078bde2c02
Injector.drop, CompletionResolver.fail, SessionManager.onFailed/reapIdle/
acquire spawn failures, and MessageService.abandon used to fail a member
or a caller's request with no log, a bare DEBUG, or a log that named only
the symptom ("session marked failed", "failed send via turn-stall
fallback"). Each now logs at WARN and names the real cause and the
numbers involved. Observability only — no behaviour changed.
570 lines
28 KiB
Java
570 lines
28 KiB
Java
package dev.ltms.bridged.msg;
|
|
|
|
import dev.ltms.bridged.herdr.AgentControl;
|
|
import dev.ltms.bridged.herdr.AgentStatus;
|
|
import dev.ltms.bridged.inject.Injector;
|
|
import dev.ltms.bridged.metrics.BridgedMetrics;
|
|
import dev.ltms.bridged.metrics.Metrics;
|
|
import org.slf4j.Logger;
|
|
import org.slf4j.LoggerFactory;
|
|
|
|
import java.util.List;
|
|
import java.util.UUID;
|
|
import java.util.concurrent.CompletableFuture;
|
|
import java.util.concurrent.CompletionException;
|
|
import java.util.concurrent.ConcurrentHashMap;
|
|
import java.util.concurrent.ExecutionException;
|
|
import java.util.concurrent.ExecutorService;
|
|
import java.util.concurrent.Executors;
|
|
import java.util.concurrent.TimeUnit;
|
|
import java.util.concurrent.TimeoutException;
|
|
import java.util.concurrent.atomic.AtomicLong;
|
|
import java.util.concurrent.locks.ReentrantLock;
|
|
|
|
/**
|
|
* The blocking delegation feature (CB-104): deliver {@code content} into a worker and block until
|
|
* the worker returns a <em>structured reply</em> via {@code bridge_reply} (the {@link Rendezvous}),
|
|
* then hand that reply back. Delivery is the {@link Injector}'s job (the background poller sends it
|
|
* when the worker is injectable); this service never drives the injector or scrapes the terminal —
|
|
* completion is the worker's explicit reply, not a guess about {@code agent_status}.
|
|
*
|
|
* <p>Sends are serialized per session so exactly one reply can be outstanding per worker, which is
|
|
* what lets a reply map unambiguously to its send (no cross-talk between concurrent callers).
|
|
*
|
|
* <p>If the worker never replies within the timeout, the caller gets a typed "still working" /
|
|
* "queued" outcome — the message may still be mid-flight. A finished-but-unreplied turn is caught
|
|
* by the CB-106 completion fallback (see {@link Rendezvous#resolveCompletion}).
|
|
*
|
|
* <p><strong>Async fire-and-poll (CB-107).</strong> A caller's MCP client caps a blocking call at
|
|
* ~60s, but a real delegated task runs for minutes. {@link #sendAsync} therefore runs the same
|
|
* blocking {@link #send} on a background virtual thread and hands back a <em>ticket</em> the caller
|
|
* polls with {@link #poll}. The blocking and async paths share one code path (and the same per-target
|
|
* serialization), so async inherits the reply + completion resolution behaviour for free.
|
|
*/
|
|
public final class MessageService {
|
|
|
|
private static final Logger log = LoggerFactory.getLogger(MessageService.class);
|
|
|
|
/**
|
|
* The window a fire-and-poll send waits for resolution — generous, since no caller is blocked on
|
|
* it; a real delegated task resolves (reply or completion) well within this, and only a genuinely
|
|
* hung worker rides it out.
|
|
*/
|
|
private static final long ASYNC_TIMEOUT_MS = 30 * 60 * 1_000L;
|
|
|
|
/** How long a finished (terminal) ticket is retained for polling before it is pruned. */
|
|
private static final long TICKET_TTL_NANOS = 10 * 60 * 1_000_000_000L;
|
|
|
|
/** Outcome of a blocking send. */
|
|
public enum Outcome {
|
|
/** The worker called {@code bridge_reply}; {@code text} holds the structured answer. */
|
|
REPLIED,
|
|
/**
|
|
* The worker's delegated turn finished without a {@code bridge_reply} (CB-106 fallback);
|
|
* {@code text} is the scraped transcript tail rather than a structured answer.
|
|
*/
|
|
COMPLETED_UNREPLIED,
|
|
/**
|
|
* The worker ran the turn then wedged in an unrecoverable state (CB-109); {@code text} is the
|
|
* failure context (e.g. the error screen). Terminal, but not a successful completion.
|
|
*/
|
|
WORKER_FAILED,
|
|
/**
|
|
* The worker paused mid-turn to ask the primary a question (CB-205); {@code text} is the
|
|
* question and {@code turnId} correlates the answer. Not terminal — the primary answers with
|
|
* {@link #answer(String, String, long)} and the turn resumes.
|
|
*/
|
|
QUESTION,
|
|
/** Timed out after the message was delivered — the worker is still working. */
|
|
TIMED_OUT_WORKING,
|
|
/** Timed out before delivery — the message is still queued for the worker. */
|
|
TIMED_OUT_QUEUED,
|
|
/** Another send to this session was in flight for the whole window. */
|
|
BUSY,
|
|
/**
|
|
* An answer ({@link #answer(String, String, long)}) referenced a {@code turnId} that is no
|
|
* longer open — the worker's {@code bridge_ask} already timed out or was answered.
|
|
*/
|
|
STALE_TURN
|
|
}
|
|
|
|
/**
|
|
* @param outcome how the send ended (or paused)
|
|
* @param text the worker's answer when {@link #completed()} (a structured {@code bridge_reply}
|
|
* for {@link Outcome#REPLIED}, a scraped transcript tail for
|
|
* {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION},
|
|
* else {@code null}
|
|
* @param turnId correlation id for a {@link Outcome#QUESTION} (answered via
|
|
* {@link #answer(String, String, long)}), else {@code null}
|
|
*/
|
|
public record Reply(Outcome outcome, String text, String turnId) {
|
|
/** A reply with no correlation id (the common terminal outcomes). */
|
|
public Reply(Outcome outcome, String text) {
|
|
this(outcome, text, null);
|
|
}
|
|
|
|
/** Whether the worker's turn actually finished with an answer (replied or scraped). */
|
|
public boolean completed() {
|
|
return outcome == Outcome.REPLIED || outcome == Outcome.COMPLETED_UNREPLIED;
|
|
}
|
|
}
|
|
|
|
/** How a worker's {@code bridge_ask} (CB-205) resolved. */
|
|
public enum AskOutcome {
|
|
/** The primary answered; {@link AskResult#answer} carries it. */
|
|
ANSWERED,
|
|
/** No delegation was open to surface the question to — the worker has no one to ask. */
|
|
NO_WAITER,
|
|
/** The primary did not answer within the window. */
|
|
TIMED_OUT
|
|
}
|
|
|
|
/** The outcome of a worker's {@code bridge_ask}: how it resolved and (if answered) the answer. */
|
|
public record AskResult(AskOutcome outcome, String answer) {
|
|
}
|
|
|
|
/** Lifecycle phase of an async delegation ticket. */
|
|
public enum Phase {
|
|
/** Delegated and in flight — queued for the worker or being worked. */
|
|
PENDING,
|
|
/** The worker's turn finished; {@link TaskView#reply} holds the answer. */
|
|
DONE,
|
|
/** The delegation could not complete (timed out, worker gone, or busy). */
|
|
FAILED
|
|
}
|
|
|
|
/**
|
|
* A poll snapshot of an async delegation.
|
|
*
|
|
* @param reply the answer when {@link #phase} is {@link Phase#DONE}, else {@code null}
|
|
* @param replySource {@code "reply"} (structured {@code bridge_reply}) or {@code "transcript"}
|
|
* (completion scrape) when {@link Phase#DONE}, else {@code null}
|
|
* @param detail a human note (live worker status while pending, or the failure reason)
|
|
*/
|
|
public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail) {
|
|
}
|
|
|
|
/** An in-flight or finished async delegation, keyed by its ticket. */
|
|
private record Task(String target, CompletableFuture<Reply> future, long createdNanos) {
|
|
}
|
|
|
|
private final AgentControl agents;
|
|
private final Injector injector;
|
|
private final Rendezvous rendezvous;
|
|
private final ReplyInbox inbox;
|
|
private final ReplyPushLoop pushLoop;
|
|
private final Metrics metrics; // CB-502: nullable — no registry in unit tests
|
|
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
|
|
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
|
|
private final AtomicLong ticketSeq = new AtomicLong();
|
|
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
|
|
Thread.ofVirtual().name("bridge-async-", 0).factory());
|
|
|
|
/**
|
|
* Create with an explicit {@link ReplyInbox} and optional {@link ReplyPushLoop}.
|
|
*
|
|
* @param pushLoop nullable — when non-null, the push loop is notified on the no-waiter reply
|
|
* branch ({@link #reply}) so it can nudge the primary to drain the inbox
|
|
*/
|
|
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
|
|
ReplyInbox inbox, ReplyPushLoop pushLoop) {
|
|
this(agents, injector, rendezvous, inbox, pushLoop, null);
|
|
}
|
|
|
|
/**
|
|
* As above, with a metric registry (CB-502). Instrumenting here rather than at the REST and MCP
|
|
* edges means both surfaces are counted by one piece of code and cannot drift.
|
|
*
|
|
* @param metrics nullable — when null, nothing is recorded
|
|
*/
|
|
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
|
|
ReplyInbox inbox, ReplyPushLoop pushLoop, Metrics metrics) {
|
|
this.agents = agents;
|
|
this.injector = injector;
|
|
this.rendezvous = rendezvous;
|
|
this.inbox = inbox;
|
|
this.pushLoop = pushLoop;
|
|
this.metrics = metrics;
|
|
}
|
|
|
|
/** Create with an explicit {@link ReplyInbox} and no push loop. */
|
|
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox) {
|
|
this(agents, injector, rendezvous, inbox, null);
|
|
}
|
|
|
|
/** Backward-compatible constructor that uses a default {@link InMemoryReplyInbox}. */
|
|
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous) {
|
|
this(agents, injector, rendezvous, new InMemoryReplyInbox());
|
|
}
|
|
|
|
/** Current lifecycle status of a worker (the {@code GET /sessions/{id}/status} surface). */
|
|
public AgentStatus status(String target) {
|
|
return agents.status(target);
|
|
}
|
|
|
|
/**
|
|
* Route a worker's explicit {@code bridge_reply}: resolve an open send, or queue it in the
|
|
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
|
|
* result is <em>not</em> a failure — the reply is held for later drain.
|
|
*
|
|
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code bridge_ask} /
|
|
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
|
* are interactive and must never be queued.
|
|
*
|
|
* @return always {@code true} — the reply either resolved a live send or was queued
|
|
*/
|
|
public boolean reply(String session, String content) {
|
|
if (rendezvous.resolve(session, content)) {
|
|
count(BridgedMetrics.REPLIES, "path", "rendezvous");
|
|
return true; // a live send took it — unchanged fast path
|
|
}
|
|
inbox.publish(session, UUID.randomUUID().toString(), content);
|
|
// A rising inbox share is the signal CB-307 exists to make visible: the worker finished but
|
|
// nobody was waiting, so delivery now depends on the push loop and a drain.
|
|
count(BridgedMetrics.REPLIES, "path", "inbox");
|
|
if (pushLoop != null) {
|
|
pushLoop.onReplyQueued(session);
|
|
}
|
|
return true; // held, not lost
|
|
}
|
|
|
|
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
|
|
private void count(String name, String... labels) {
|
|
if (metrics != null) {
|
|
metrics.inc(name, labels);
|
|
}
|
|
}
|
|
|
|
/** Count a send's terminal outcome and pass the reply through unchanged. */
|
|
private Reply recorded(Reply r) {
|
|
String label = sendOutcomeLabel(r.outcome());
|
|
if (label != null) {
|
|
count(BridgedMetrics.SENDS, "outcome", label);
|
|
}
|
|
return r;
|
|
}
|
|
|
|
/** Map a terminal send outcome to its metric label, or {@code null} for non-terminal ones. */
|
|
private static String sendOutcomeLabel(Outcome o) {
|
|
return switch (o) {
|
|
case REPLIED -> "replied";
|
|
case COMPLETED_UNREPLIED -> "completion_fallback";
|
|
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
|
|
case WORKER_FAILED -> "failed";
|
|
case STALE_TURN, QUESTION -> null; // not a completed delegation
|
|
};
|
|
}
|
|
|
|
/**
|
|
* Abandon any send still waiting on {@code target} because its session has gone away (CB-516).
|
|
*
|
|
* <p>Without this, tearing a worker down left its rendezvous waiter open: a blocking
|
|
* {@code bridge_send} kept blocking, and an async one kept reporting {@code PENDING} until
|
|
* {@link #ASYNC_TIMEOUT_MS} — thirty minutes — even though the worker provably no longer
|
|
* existed and the delegation could never complete. Worse, {@code poll} already had the evidence
|
|
* (it calls {@code liveStatus} to build its detail string and gets back {@code "unknown"}) and
|
|
* reported {@code PENDING} anyway.
|
|
*
|
|
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
|
|
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
|
|
*
|
|
* @return true if a live waiter was failed
|
|
*/
|
|
public boolean abandon(String target, String reason) {
|
|
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
|
if (waiter == null || waiter.isDone()) {
|
|
return false; // nobody is blocked on this worker — nothing to abandon
|
|
}
|
|
boolean failed = rendezvous.resolveFailure(waiter, reason);
|
|
if (failed) {
|
|
log.warn("abandoning the blocked send to {}: {}", target, reason);
|
|
}
|
|
return failed;
|
|
}
|
|
|
|
/**
|
|
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
|
* so that a subsequent drain or peek no longer returns it.
|
|
*/
|
|
public void ackReply(String target, String msgId) {
|
|
inbox.ack(target, msgId);
|
|
}
|
|
|
|
/**
|
|
* Drain (peek + ack) all pending inbox replies for {@code target}. At-least-once: returns the
|
|
* messages and acknowledges them; an in-flight failure between returning and the caller
|
|
* processing them re-surfaces them on a subsequent drain (the ack is local).
|
|
*
|
|
* @return the drained messages, newest last (FIFO); empty list if none
|
|
*/
|
|
public List<ReplyInbox.InboxMessage> drainReplies(String target) {
|
|
var messages = inbox.peek(target);
|
|
for (var msg : messages) {
|
|
inbox.ack(target, msg.msgId());
|
|
}
|
|
return messages;
|
|
}
|
|
|
|
/**
|
|
* Deliver {@code content} to {@code target} (a herdr {@code terminal_id}) and block until the
|
|
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses.
|
|
*/
|
|
public Reply send(String target, String content, long timeoutMillis) {
|
|
return send(target, content, timeoutMillis, null);
|
|
}
|
|
|
|
/**
|
|
* As {@link #send(String, String, long)}, but with an accepted-delivery hook.
|
|
*
|
|
* <p>{@code onAccepted} is invoked exactly once, once this send has won {@code target}'s send
|
|
* lock and so become the <em>accepted target turn</em> — it runs <em>before</em> delivery is
|
|
* queued, so a throwing hook fails the send cleanly (the waiter it already opened is closed and
|
|
* nothing is left queued). It is <em>not</em> invoked when the send is {@link Outcome#BUSY}
|
|
* (lock never taken). A caller uses this to record that <em>it</em> now owns the delegation's
|
|
* reply routing (CB-548: {@code PrimaryRegistry} delegator ownership) — recording only on
|
|
* acceptance means a concurrent sender that times out {@code BUSY} can never steal ownership it
|
|
* never earned. {@code null} disables the hook.
|
|
*/
|
|
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted) {
|
|
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
|
|
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
|
|
|
|
if (!tryLock(lock, remainingMillis(deadlineNanos))) {
|
|
return new Reply(Outcome.BUSY, null); // another send held the session the whole window
|
|
}
|
|
try {
|
|
// Open the waiter BEFORE queueing delivery (CB-548). A fast reply — the worker already
|
|
// injectable the instant we enqueue — otherwise arrives before the waiter is registered
|
|
// and orphans into the inbox while this send blocks to the timeout (the enqueue-before-
|
|
// open race). Opening first also means a throwing onAccepted (fired before enqueue) or an
|
|
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
|
|
// failed send leaves no stale waiter behind.
|
|
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
|
|
try {
|
|
// The send has won the lock; the accepted-delivery hook records delegator ownership
|
|
// here (CB-548). It runs BEFORE enqueue so a throwing hook — onAccepted is now a
|
|
// public callback — fails the send without queuing a message that would orphan.
|
|
if (onAccepted != null) {
|
|
onAccepted.run();
|
|
}
|
|
CompletableFuture<Void> delivered = injector.enqueue(target, content);
|
|
try {
|
|
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
|
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
|
|
} catch (TimeoutException e) {
|
|
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
|
|
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
|
|
return recorded(new Reply(
|
|
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
|
|
} catch (ExecutionException e) {
|
|
Throwable cause = e.getCause();
|
|
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
|
} catch (InterruptedException e) {
|
|
Thread.currentThread().interrupt();
|
|
throw new IllegalStateException("interrupted awaiting reply from " + target, e);
|
|
}
|
|
} finally {
|
|
rendezvous.close(target, reply);
|
|
}
|
|
} finally {
|
|
lock.unlock();
|
|
}
|
|
}
|
|
|
|
/**
|
|
* A worker's mid-turn question (CB-205 reverse rendezvous): surface {@code question} to the
|
|
* primary by resolving its open blocking {@code bridge_send}, then block this (worker) call until
|
|
* the primary answers via {@link #answer} or {@code timeoutMillis} elapses. Identity is the
|
|
* worker's own session — it does not address the primary.
|
|
*
|
|
* <p>Returns {@link AskOutcome#NO_WAITER} when no delegation is open to surface the question to
|
|
* (nothing to answer it), {@link AskOutcome#ANSWERED} with the primary's answer, or
|
|
* {@link AskOutcome#TIMED_OUT} if the primary stayed silent. The worker resumes its turn either
|
|
* way — an answered ask hands back the answer; an unanswered one leaves it to proceed alone.
|
|
*/
|
|
public AskResult ask(String workerSession, String question, long timeoutMillis) {
|
|
Rendezvous.AskTicket ticket = rendezvous.openAsk(workerSession);
|
|
// Only the freshly-opening caller surfaces the question; a coalesced duplicate simply blocks on
|
|
// the shared answer future that the fresh owner is already responsible for.
|
|
if (ticket.fresh()) {
|
|
// Register the reverse waiter first, then surface the question — so the answer, which can
|
|
// arrive the instant the primary reacts, always finds an open waiter to resolve.
|
|
if (!rendezvous.resolveQuestion(workerSession, question, ticket.turnId())) {
|
|
rendezvous.closeAsk(ticket.turnId());
|
|
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
|
|
}
|
|
}
|
|
try {
|
|
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
|
|
return new AskResult(AskOutcome.ANSWERED, answer);
|
|
} catch (TimeoutException e) {
|
|
log.debug("bridge_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
|
|
return new AskResult(AskOutcome.TIMED_OUT, null);
|
|
} catch (ExecutionException e) {
|
|
Throwable cause = e.getCause();
|
|
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
|
} catch (InterruptedException e) {
|
|
Thread.currentThread().interrupt();
|
|
throw new IllegalStateException("interrupted awaiting the primary's answer for " + workerSession, e);
|
|
} finally {
|
|
// Only the fresh owner tears down the shared turn; a duplicate must leave it open.
|
|
if (ticket.fresh()) {
|
|
rendezvous.closeAsk(ticket.turnId());
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* The primary's answer to a worker's {@code bridge_ask} (CB-205): resolve the worker's blocked
|
|
* question identified by {@code turnId}, then — like a fresh {@link #send} — block for the worker's
|
|
* eventual {@code bridge_reply} as it finishes the resumed turn. The worker session is derived from
|
|
* {@code turnId}, never a caller argument.
|
|
*
|
|
* <p>Unlike {@link #send} this does not re-inject through the {@link Injector}: the worker is
|
|
* mid-turn (already picked up), so the answer flows back through its own open {@code bridge_ask}
|
|
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
|
|
* is unblocked so a reply that lands the instant it resumes is not lost.
|
|
*/
|
|
public Reply answer(String turnId, String content, long timeoutMillis) {
|
|
String workerSession = rendezvous.askSession(turnId);
|
|
if (workerSession == null) {
|
|
return new Reply(Outcome.STALE_TURN, null); // the ask lapsed (timed out or already answered)
|
|
}
|
|
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
|
|
ReentrantLock lock = sessionLocks.computeIfAbsent(workerSession, _ -> new ReentrantLock());
|
|
if (!tryLock(lock, remainingMillis(deadlineNanos))) {
|
|
return new Reply(Outcome.BUSY, null);
|
|
}
|
|
try {
|
|
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
|
|
if (!rendezvous.answerAsk(turnId, content)) {
|
|
rendezvous.close(workerSession, reply);
|
|
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
|
|
}
|
|
try {
|
|
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
|
return new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
|
|
} catch (TimeoutException e) {
|
|
// The worker resumed but hasn't replied yet — no completion fallback arms an answered
|
|
// turn (it never re-entered the injector), so a silent worker rides out the window.
|
|
return new Reply(Outcome.TIMED_OUT_WORKING, null);
|
|
} catch (ExecutionException e) {
|
|
Throwable cause = e.getCause();
|
|
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
|
} catch (InterruptedException e) {
|
|
Thread.currentThread().interrupt();
|
|
throw new IllegalStateException("interrupted awaiting reply from " + workerSession, e);
|
|
} finally {
|
|
rendezvous.close(workerSession, reply);
|
|
}
|
|
} finally {
|
|
lock.unlock();
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Fire-and-poll variant of {@link #send}: deliver {@code content} to {@code target} on a
|
|
* background virtual thread and return immediately with a ticket to {@link #poll}. This is how a
|
|
* long task is delegated without tripping the caller's MCP client call timeout.
|
|
*
|
|
* @return the ticket to poll for the eventual result
|
|
*/
|
|
public String sendAsync(String target, String content) {
|
|
return sendAsync(target, content, null);
|
|
}
|
|
|
|
/**
|
|
* As {@link #sendAsync(String, String)}, with the accepted-delivery hook of
|
|
* {@link #send(String, String, long, Runnable)} — the running {@code send} invokes {@code onAccepted}
|
|
* the moment it becomes the accepted target turn, so async flooding records delegator ownership
|
|
* exactly as the blocking path does (CB-548).
|
|
*
|
|
* @return the ticket to poll for the eventual result
|
|
*/
|
|
public String sendAsync(String target, String content, Runnable onAccepted) {
|
|
String ticket = "task-" + ticketSeq.incrementAndGet();
|
|
CompletableFuture<Reply> future = CompletableFuture.supplyAsync(
|
|
() -> send(target, content, ASYNC_TIMEOUT_MS, onAccepted), asyncExecutor);
|
|
tasks.put(ticket, new Task(target, future, System.nanoTime()));
|
|
pruneTerminalTickets();
|
|
log.debug("async send {} -> {}", ticket, target);
|
|
return ticket;
|
|
}
|
|
|
|
/**
|
|
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket;
|
|
* otherwise a {@link Phase#PENDING} view (with the live worker status as detail), a
|
|
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
|
|
*/
|
|
public TaskView poll(String ticket) {
|
|
Task task = tasks.get(ticket);
|
|
if (task == null) {
|
|
return null;
|
|
}
|
|
CompletableFuture<Reply> f = task.future();
|
|
if (!f.isDone()) {
|
|
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target()));
|
|
}
|
|
Reply r;
|
|
try {
|
|
r = f.getNow(null);
|
|
} catch (CompletionException | java.util.concurrent.CancellationException e) {
|
|
Throwable cause = (e instanceof CompletionException ce && ce.getCause() != null) ? ce.getCause() : e;
|
|
return new TaskView(ticket, Phase.FAILED, null, null, cause.getMessage());
|
|
}
|
|
if (r.completed()) {
|
|
String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript";
|
|
return new TaskView(ticket, Phase.DONE, r.text(), source, null);
|
|
}
|
|
// A wedged worker (CB-109) carries the error context as its reason; the timeout/busy
|
|
// outcomes carry none, so fall back to the outcome name.
|
|
String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null
|
|
? r.text()
|
|
: "no reply — " + r.outcome().name().toLowerCase();
|
|
return new TaskView(ticket, Phase.FAILED, null, null, detail);
|
|
}
|
|
|
|
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
|
|
private String liveStatus(String target) {
|
|
try {
|
|
return agents.status(target).name().toLowerCase();
|
|
} catch (RuntimeException e) {
|
|
return "unknown";
|
|
}
|
|
}
|
|
|
|
/** Drop finished tickets older than the TTL so the registry cannot grow without bound. */
|
|
private void pruneTerminalTickets() {
|
|
long cutoff = System.nanoTime() - TICKET_TTL_NANOS;
|
|
tasks.values().removeIf(t -> t.future().isDone() && t.createdNanos() < cutoff);
|
|
}
|
|
|
|
/** Release the async executor. */
|
|
public void close() {
|
|
asyncExecutor.shutdown();
|
|
}
|
|
|
|
/** Map a rendezvous {@link Rendezvous.Kind} onto its send {@link Outcome} (shared by send/answer). */
|
|
private static Outcome outcomeOf(Rendezvous.Kind kind) {
|
|
return switch (kind) {
|
|
case REPLY -> Outcome.REPLIED;
|
|
case COMPLETION -> Outcome.COMPLETED_UNREPLIED;
|
|
case FAILED -> Outcome.WORKER_FAILED;
|
|
case QUESTION -> Outcome.QUESTION;
|
|
};
|
|
}
|
|
|
|
private static boolean tryLock(ReentrantLock lock, long millis) {
|
|
try {
|
|
return lock.tryLock(Math.max(0, millis), TimeUnit.MILLISECONDS);
|
|
} catch (InterruptedException e) {
|
|
Thread.currentThread().interrupt();
|
|
throw new IllegalStateException("interrupted awaiting the session send lock", e);
|
|
}
|
|
}
|
|
|
|
private static long remainingMillis(long deadlineNanos) {
|
|
return (deadlineNanos - System.nanoTime()) / 1_000_000L;
|
|
}
|
|
}
|