9b50dd69d8
Part of #145 (CB-632), under epic #125. The product is called fleet and the daemon is called fleetd, but the code still said bridge everywhere. This renames the Java half: package dev.ltms.bridged -> dev.ltms.fleet Bridged -> Fleetd (the main class) BridgedConfig -> FleetConfig BridgeMcp -> FleetMcp BridgedApp -> FleetApp BridgedMetrics -> FleetMetrics The package root is dev.ltms.fleet, not dev.ltms.fleetd. The trailing d means daemon, which names a process, not a namespace. What this commit deliberately does NOT change: - The module directory stays bridged/, and <finalName> stays bridged. The installed launchd plist names bridged/target/bridged.jar and its KeepAlive is armed, so renaming the jar on its own strands a restart. Both change at the cutover, together with the plist, in one step. - The bridge_* MCP tool aliases. CB-622 shipped both names on purpose. One test names a local variable viaBridge because it holds the result of the deprecated call; the rename collided with it and the compiler caught it. That variable is back. - BRIDGED_* env var names, and bridged.yaml. Both are operator contracts and need a read-both shim, which is a later unit. Two things a plain search-and-replace would have missed: - logback.xml and logback-test.xml name the package twice, once as a turboFilter class= attribute. The compiler never checks those. - BSD sed does not support \b. The word-boundary expression matched nothing and said nothing, while the other ten in the same command worked. Checked the leftovers instead of trusting the exit code. Verified: mvn clean install green, 51 test classes, 878 tests, 0 failures -- the same count as before the rename.
832 lines
42 KiB
Java
832 lines
42 KiB
Java
package dev.ltms.fleet.msg;
|
|
|
|
import dev.ltms.fleet.herdr.AgentControl;
|
|
import dev.ltms.fleet.herdr.AgentStatus;
|
|
import dev.ltms.fleet.inject.Injector;
|
|
import dev.ltms.fleet.metrics.FleetMetrics;
|
|
import dev.ltms.fleet.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;
|
|
import java.util.function.LongSupplier;
|
|
|
|
/**
|
|
* 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 fleet_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. Package-
|
|
* private (not {@code private}) so a test can advance an injected clock past it deterministically
|
|
* instead of duplicating the magic number or sleeping for real.
|
|
*/
|
|
static final long TICKET_TTL_NANOS = 10 * 60 * 1_000_000_000L;
|
|
|
|
/** Outcome of a blocking send. */
|
|
public enum Outcome {
|
|
/** The worker called {@code fleet_reply}; {@code text} holds the structured answer. */
|
|
REPLIED,
|
|
/**
|
|
* The worker's delegated turn finished without a {@code fleet_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 turn finished without a {@code fleet_reply} and the scrape matched the backend's
|
|
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
|
|
* carrying the matched line. The worker's pane is healthy — only its account is refusing —
|
|
* so this is never reported as a completed reply, and is kept distinct from
|
|
* {@link #WORKER_FAILED} (a wedged worker) and a session simply going {@code GONE}.
|
|
*/
|
|
BACKEND_EXHAUSTED,
|
|
/**
|
|
* 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 fleet_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 fleet_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 fleet_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 fleet_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 is paused in {@code fleet_ask}; {@link TaskView#reply} and {@link TaskView#turnId} identify it. */
|
|
ASKING,
|
|
/** 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}, or the question when
|
|
* {@link #phase} is {@link Phase#ASKING}; otherwise {@code null}
|
|
* @param replySource {@code "reply"} (structured {@code fleet_reply}) or {@code "transcript"}
|
|
* (completion scrape) when {@link Phase#DONE}, else {@code null}
|
|
* @param detail a human note (live worker status while pending, ask state, or failure reason)
|
|
* @param turnId correlation id for an {@link Phase#ASKING} ticket, else {@code null}
|
|
*/
|
|
public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail,
|
|
String turnId) {
|
|
}
|
|
|
|
/** An in-flight or finished async delegation, keyed by its ticket. */
|
|
private static final class Task {
|
|
private final String ticket;
|
|
private final String target;
|
|
private final CompletableFuture<Reply> future = new CompletableFuture<>();
|
|
private final long createdNanos;
|
|
private volatile Reply question;
|
|
private volatile String turnId;
|
|
|
|
private Task(String ticket, String target, long createdNanos) {
|
|
this.ticket = ticket;
|
|
this.target = target;
|
|
this.createdNanos = createdNanos;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* A worker session's currently-open {@code fleet_ask} question, surfaced so {@code fleet_status}
|
|
* can show it without the caller needing the ticket first (CB-582). Only covers async
|
|
* (fire-and-poll) delegations, which track the question on their {@link Task}; a blocking
|
|
* ({@code wait:true}) send already hands the question straight back to its own caller, so there is
|
|
* nothing hidden left for {@code fleet_status} to surface in that case.
|
|
*/
|
|
public record PendingAsk(String ticket, String question, String turnId) {
|
|
}
|
|
|
|
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
|
|
// CB-588: injectable so pruneTerminalTickets' 10-minute TICKET_TTL_NANOS can be exercised in a
|
|
// test without a real wait — same seam SessionManager already uses for its idle reaper (nowNanos).
|
|
private final LongSupplier nowNanos;
|
|
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
|
|
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
|
|
/** Async task that owns each exact forward rendezvous waiter. */
|
|
private final ConcurrentHashMap<CompletableFuture<Rendezvous.Resolution>, Task> asyncTasksByWaiter =
|
|
new ConcurrentHashMap<>();
|
|
/** Async tickets paused on a specific {@code fleet_ask} turn. */
|
|
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = 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,
|
|
* (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a
|
|
* terminal phase, whenever {@link #poll} hands a terminal ticket to its caller,
|
|
* and (CB-582) whenever an async ticket's worker pauses mid-turn in
|
|
* {@code fleet_ask} or that pause ends (answered or lapsed)
|
|
*/
|
|
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, injector, rendezvous, inbox, pushLoop, metrics, System::nanoTime);
|
|
}
|
|
|
|
/** Test constructor with an injectable clock (CB-588: exercise the ticket-prune TTL without a real wait). */
|
|
MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox,
|
|
ReplyPushLoop pushLoop, Metrics metrics, LongSupplier nowNanos) {
|
|
this.agents = agents;
|
|
this.injector = injector;
|
|
this.rendezvous = rendezvous;
|
|
this.inbox = inbox;
|
|
this.pushLoop = pushLoop;
|
|
this.metrics = metrics;
|
|
this.nowNanos = nowNanos;
|
|
}
|
|
|
|
/** 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);
|
|
}
|
|
|
|
/** Read-only delegation fact for fleet views. */
|
|
public boolean hasAcceptedDelivery(String target) {
|
|
return rendezvous.isWaiting(target);
|
|
}
|
|
|
|
/** Read-only inbox fact for fleet views. */
|
|
public boolean hasInboxMessage(String target) {
|
|
return !inbox.peek(target).isEmpty();
|
|
}
|
|
|
|
/**
|
|
* Route a worker's explicit {@code fleet_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 fleet_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(FleetMetrics.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(FleetMetrics.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(FleetMetrics.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 BACKEND_EXHAUSTED -> "backend_exhausted";
|
|
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 fleet_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);
|
|
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
|
boolean asyncFailed = false;
|
|
for (Task task : tasks.values()) {
|
|
if (target.equals(task.target) && task.question == null
|
|
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
|
|
asyncFailed = true;
|
|
}
|
|
}
|
|
if (failed) {
|
|
log.warn("abandoning the blocked send to {}: {}", target, reason);
|
|
}
|
|
return failed || asyncFailed;
|
|
}
|
|
|
|
/**
|
|
* 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}.
|
|
*
|
|
* <p><strong>The ack happens here, before the caller has the messages</strong> — before the MCP
|
|
* or REST response carrying them has been written, and long before the client has processed
|
|
* them. That ordering is what the two adapters disagree about, so do not read this method as
|
|
* "at-least-once" without qualifying which inbox is behind it (CB-529):
|
|
*
|
|
* <ul>
|
|
* <li>{@code InMemoryReplyInbox} — the ack only drops an entry from a local map. The messages
|
|
* are already in the returned list, so nothing can be lost after this point.
|
|
* <li>{@code AmqpReplyInbox} — the ack is a broker-side {@code basicAck}. Once it lands the
|
|
* broker has forgotten the message. If the daemon dies while writing the response, the
|
|
* reply is gone from the broker <em>and</em> the client never received it. Re-polling
|
|
* cannot recover it, because there is nothing left to re-deliver.
|
|
* </ul>
|
|
*
|
|
* <p>So the loss window is the response write, and it is a genuine loss rather than a
|
|
* redelivery. This is accepted, not overlooked: the alternative — ack on the next poll — turns
|
|
* every normal drain into a double delivery, which costs more than the window it closes. A
|
|
* caller that needs certainty re-polls; that is idempotent for every case except this one.
|
|
*
|
|
* <p>Any change here must be checked against <em>both</em> adapters. The previous version of
|
|
* this javadoc claimed "the ack is local", which was true when only the in-memory inbox existed
|
|
* and silently became false when the AMQP adapter landed.
|
|
*
|
|
* @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) {
|
|
return send(target, content, timeoutMillis, onAccepted, null);
|
|
}
|
|
|
|
/** Run a send, optionally stopping an async task that teardown already failed before acceptance. */
|
|
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task) {
|
|
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 {
|
|
if (task != null && task.future.isDone()) {
|
|
return task.future.getNow(null);
|
|
}
|
|
if (hasAsyncQuestion(target)) {
|
|
return new Reply(Outcome.BUSY, null); // the worker's current turn is paused for its lead
|
|
}
|
|
// 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 {
|
|
if (task != null) {
|
|
asyncTasksByWaiter.put(reply, task);
|
|
}
|
|
TurnToken token = new TurnToken(target, reply);
|
|
// 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, token);
|
|
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 {
|
|
asyncTasksByWaiter.remove(reply);
|
|
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 fleet_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.
|
|
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(workerSession);
|
|
Task task = markAsyncQuestion(waiter, question, ticket.turnId());
|
|
if (!rendezvous.resolveQuestion(workerSession, question, ticket.turnId())) {
|
|
if (task != null) {
|
|
clearAsyncQuestion(ticket.turnId(), true);
|
|
}
|
|
rendezvous.closeAsk(ticket.turnId());
|
|
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
|
|
}
|
|
// CB-582: the question just became visible via fleet_poll (Phase.ASKING) for an async
|
|
// (wait:false) delegation — nudge the lead's own pane the same way a terminal ticket does
|
|
// (CB-588), since the lead's normal poll cadence is minutes away and the reverse-rendezvous
|
|
// window (~55s, see FleetMcp/FleetApp) is far shorter. A blocking (wait:true) send has
|
|
// no Task and gets the question directly in its own reply, so task == null there — nothing
|
|
// to nudge.
|
|
if (task != null && pushLoop != null) {
|
|
pushLoop.onQuestionOpened(task.ticket, workerSession, ticket.turnId(), question);
|
|
}
|
|
}
|
|
try {
|
|
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
|
|
return new AskResult(AskOutcome.ANSWERED, answer);
|
|
} catch (TimeoutException e) {
|
|
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
|
|
clearAsyncQuestion(ticket.turnId(), true);
|
|
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());
|
|
// CB-582: tear the push loop's copy down at the same point, not only on the three
|
|
// paths that call clearAsyncQuestion. The answer future can complete exceptionally
|
|
// (ExecutionException) or the thread be interrupted, and both leave this method by
|
|
// throwing — the question would stay pending forever, keep being named in nudges
|
|
// until its own cap, and never be removed from the map. Already-closed is a no-op.
|
|
if (pushLoop != null) {
|
|
pushLoop.questionClosed(ticket.turnId());
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* The primary's answer to a worker's {@code fleet_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 fleet_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 fleet_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
|
|
}
|
|
clearAsyncQuestion(turnId, false);
|
|
try {
|
|
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
|
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
|
|
finishAsyncTask(turnId, result);
|
|
return result;
|
|
} 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();
|
|
Task task = new Task(ticket, target, nowNanos.getAsLong());
|
|
tasks.put(ticket, task);
|
|
if (pushLoop != null) {
|
|
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
|
// worker paused in fleet_ask leaves it running, per finishAsyncTask's own contract — so
|
|
// this fires exactly once, from whichever path completes it: finishAsyncTask(task, result)
|
|
// below on any non-QUESTION outcome of send() — a worker's fleet_reply, the CB-106
|
|
// completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same
|
|
// finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION
|
|
// is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516
|
|
// abandon() on teardown. Without this, MessageService.reply's rendezvous fast path (the
|
|
// one an async ticket always takes) never told the push loop anything happened — see the
|
|
// class javadoc on sendAsync/CB-107.
|
|
task.future.whenComplete((reply, ex) -> {
|
|
boolean failed = ex != null || reply == null || !reply.completed();
|
|
pushLoop.onTicketTerminal(ticket, target, failed);
|
|
});
|
|
}
|
|
asyncExecutor.submit(() -> {
|
|
try {
|
|
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
|
|
if (result.outcome() == Outcome.QUESTION) {
|
|
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
|
|
// just after resolveQuestion wakes this thread.
|
|
} else {
|
|
finishAsyncTask(task, result);
|
|
}
|
|
} catch (Throwable t) {
|
|
task.future.completeExceptionally(t);
|
|
}
|
|
});
|
|
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()) {
|
|
Reply question = task.question;
|
|
if (question != null) {
|
|
return new TaskView(ticket, Phase.ASKING, question.text(), null,
|
|
"worker is waiting for your answer", question.turnId());
|
|
}
|
|
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target), null);
|
|
}
|
|
// CB-588: the ticket is terminal and being handed to the caller right here — tell the push
|
|
// loop it is collected so a later tick's nudge never names a ticket the lead already has.
|
|
if (pushLoop != null) {
|
|
pushLoop.ticketCollected(ticket);
|
|
}
|
|
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(), null);
|
|
}
|
|
if (r.completed()) {
|
|
String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript";
|
|
return new TaskView(ticket, Phase.DONE, r.text(), source, null, null);
|
|
}
|
|
// A wedged worker (CB-109) or a backend-exhausted classification (CB-578 stage A) carries
|
|
// the real cause as its reason; the timeout/busy outcomes carry none, so fall back to the
|
|
// outcome name.
|
|
boolean carriesReason = r.outcome() == Outcome.WORKER_FAILED || r.outcome() == Outcome.BACKEND_EXHAUSTED;
|
|
String detail = carriesReason && r.text() != null
|
|
? r.text()
|
|
: "no reply — " + r.outcome().name().toLowerCase();
|
|
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
|
|
}
|
|
|
|
/** 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 {@link #tasks} cannot grow without bound.
|
|
*
|
|
* <p>{@code tasks} is the sole authority on whether a ticket still exists — {@link #poll} returns
|
|
* {@code null} the instant a ticket is gone from here, before it ever reaches the terminal branch
|
|
* that calls {@link ReplyPushLoop#ticketCollected}. Without telling the push loop about a prune
|
|
* too, its own {@code pendingTickets} entry would outlive the ticket it names: an unpolled ticket
|
|
* (or one the reminder cap already gave up on) is pruned here but never collected there, so it
|
|
* lingers in {@code pendingTickets} forever and rides along on every later nudge to the same lead
|
|
* — naming a ticket {@code fleet_poll} can no longer find (CB-588 follow-up).
|
|
*/
|
|
private void pruneTerminalTickets() {
|
|
long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS;
|
|
tasks.entrySet().removeIf(e -> {
|
|
Task t = e.getValue();
|
|
boolean expired = t.future.isDone() && t.createdNanos < cutoff;
|
|
if (expired && pushLoop != null) {
|
|
pushLoop.ticketCollected(e.getKey());
|
|
}
|
|
return expired;
|
|
});
|
|
}
|
|
|
|
/** Record the active question for an async ticket; blocking sends have no entry and stay unchanged. */
|
|
private Task markAsyncQuestion(CompletableFuture<Rendezvous.Resolution> waiter, String text, String turnId) {
|
|
Task task = waiter == null ? null : asyncTasksByWaiter.get(waiter);
|
|
if (task != null) {
|
|
task.question = new Reply(Outcome.QUESTION, text, turnId);
|
|
task.turnId = turnId;
|
|
asyncTasksByTurn.put(turnId, task);
|
|
}
|
|
return task;
|
|
}
|
|
|
|
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
|
|
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
|
|
// CB-582: tell the push loop first — like ticketCollected, a removal for a turnId it never
|
|
// nudged about (or already dropped) is a harmless no-op, so this is safe to call unconditionally
|
|
// rather than threading the guard below through it.
|
|
if (pushLoop != null) {
|
|
pushLoop.questionClosed(turnId);
|
|
}
|
|
Task task = asyncTasksByTurn.get(turnId);
|
|
if (task != null && turnId.equals(task.turnId)) {
|
|
task.question = null;
|
|
if (forgetTurn) {
|
|
asyncTasksByTurn.remove(turnId, task);
|
|
task.turnId = null;
|
|
}
|
|
}
|
|
}
|
|
|
|
/** Complete and detach an async ticket after its worker's actual terminal reply. */
|
|
private void finishAsyncTask(Task task, Reply result) {
|
|
task.future.complete(result);
|
|
if (task.turnId != null) {
|
|
asyncTasksByTurn.remove(task.turnId, task);
|
|
}
|
|
}
|
|
|
|
/** Complete the async ticket correlated to a specific answered turn. */
|
|
private void finishAsyncTask(String turnId, Reply result) {
|
|
Task task = asyncTasksByTurn.get(turnId);
|
|
if (task != null) {
|
|
finishAsyncTask(task, result);
|
|
}
|
|
}
|
|
|
|
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
|
|
private boolean hasAsyncQuestion(String target) {
|
|
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
|
|
}
|
|
|
|
/**
|
|
* The question {@code workerSession} is currently paused on via {@code fleet_ask}, if any
|
|
* (CB-582) — {@code fleet_status} uses this to show a pending question without the caller
|
|
* needing the ticket. {@code null} when the session has no open async question (including a
|
|
* session mid a <em>blocking</em> {@code fleet_ask}, which has no {@link Task} to look up — see
|
|
* {@link PendingAsk}).
|
|
*/
|
|
public PendingAsk pendingAsk(String workerSession) {
|
|
for (Task task : tasks.values()) {
|
|
Reply q = task.question;
|
|
if (q != null && workerSession.equals(task.target)) {
|
|
return new PendingAsk(task.ticket, q.text(), q.turnId());
|
|
}
|
|
}
|
|
return null;
|
|
}
|
|
|
|
/** 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 BACKEND_EXHAUSTED -> Outcome.BACKEND_EXHAUSTED;
|
|
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;
|
|
}
|
|
}
|