diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index ce26254..b30c5e8 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -577,13 +577,25 @@ public final class BridgeMcp { return text("acknowledged " + msgId); } - /** {@code bridge_status}: the live lifecycle status of a worker session. */ + /** + * {@code bridge_status}: the live lifecycle status of a worker session, plus — when the worker + * is paused mid-turn in an async {@code bridge_ask} (CB-582) — the open question and how to + * answer it, so a lead on its normal poll cadence does not need the ticket to notice. + */ static McpSchema.CallToolResult status(MessageService messages, String sessionId) { if (isBlank(sessionId)) { return error("sessionId is required"); } try { - return text(messages.status(sessionId).name().toLowerCase()); + String base = messages.status(sessionId).name().toLowerCase(); + MessageService.PendingAsk ask = messages.pendingAsk(sessionId); + if (ask == null) { + return text(base); + } + return text(base + "\n\n[question — worker is waiting for your answer]\n" + ask.question() + + "\n\nAnswer it by calling bridge_send again with turnId=\"" + ask.turnId() + + "\" and content set to your answer; the worker resumes the same turn." + + " (ticket " + ask.ticket() + ")"); } catch (HerdrException e) { return error("herdr error for session " + sessionId + ": " + e.getMessage()); } diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index fe1a0b9..9926aee 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -164,18 +164,30 @@ public final class MessageService { /** 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 future = new CompletableFuture<>(); private final long createdNanos; private volatile Reply question; private volatile String turnId; - private Task(String target, long createdNanos) { + private Task(String ticket, String target, long createdNanos) { + this.ticket = ticket; this.target = target; this.createdNanos = createdNanos; } } + /** + * A worker session's currently-open {@code bridge_ask} question, surfaced so {@code bridge_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 bridge_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; @@ -200,9 +212,11 @@ public final class MessageService { * 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, and + * 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, and whenever {@link #poll} hands a terminal ticket to its caller + * 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 bridge_ask} or that pause ends (answered or lapsed) */ public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox, ReplyPushLoop pushLoop) { @@ -495,6 +509,15 @@ public final class MessageService { 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 bridge_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 BridgeMcp/BridgedApp) 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); @@ -589,7 +612,7 @@ public final class MessageService { */ public String sendAsync(String target, String content, Runnable onAccepted) { String ticket = "task-" + ticketSeq.incrementAndGet(); - Task task = new Task(target, nowNanos.getAsLong()); + 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 @@ -715,6 +738,12 @@ public final class MessageService { /** 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; @@ -746,6 +775,23 @@ public final class MessageService { return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target)); } + /** + * The question {@code workerSession} is currently paused on via {@code bridge_ask}, if any + * (CB-582) — {@code bridge_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 blocking {@code bridge_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(); diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java index 164b618..d237050 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -19,30 +19,33 @@ import java.util.stream.Collectors; /** * A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work - * waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), or an + * waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), an * async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase - * (CB-588). + * (CB-588), or an async ticket's worker pausing mid-turn in {@code bridge_ask} to await an answer + * (CB-582). * - *

CB-590: one schedule per lead. Both kinds of work are triggered through - * their own entry point — {@link #onReplyQueued(String)} and - * {@link #onTicketTerminal(String, String, boolean)} — but both resolve the lead that should be - * nudged and coalesce onto a single per-lead reminder schedule, tracked in {@link #activeLeads}. - * Earlier this was two independent schedules (one keyed by worker target for replies, one keyed - * by lead for tickets) that could both decide to inject into the same pane in the same window — - * a race, not routine behaviour, but the expensive kind: it interrupts the lead's live turn - * twice. Collapsing to one schedule per lead makes that structurally impossible: at most one - * scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}' + *

CB-590: one schedule per lead. All three kinds of work are triggered + * through their own entry point — {@link #onReplyQueued(String)}, + * {@link #onTicketTerminal(String, String, boolean)}, and + * {@link #onQuestionOpened(String, String, String, String)} — but each resolves the lead that + * should be nudged and coalesces onto a single per-lead reminder schedule, tracked in + * {@link #activeLeads}. Earlier this was two independent schedules (one keyed by worker target + * for replies, one keyed by lead for tickets) that could both decide to inject into the same pane + * in the same window — a race, not routine behaviour, but the expensive kind: it interrupts the + * lead's live turn twice. Collapsing to one schedule per lead makes that structurally impossible: + * at most one scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}' * compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight. * *

Each tick examines everything pending for that lead — reply targets whose inbox - * still holds an unacked message ({@link #pendingReplies}) and tickets not yet collected - * ({@link #pendingTickets}) — and sends at most one combined nudge per tick - * ({@link #injectNudge(String, int, int)}). Work that arrives while the lead is busy is never - * lost: it is re-read fresh on every tick until the lead is injectable or its own reminder cap - * ({@link #maxReminders}) is reached — reply and ticket work each spend from their own budget, so - * one source exhausting its cap does not stop nudges about the other (post-CB-590 regression fix; - * see {@link #decide}) — whichever the durable inbox / pending-ticket set doesn't already answer - * via {@code STOP}. + * still holds an unacked message ({@link #pendingReplies}), tickets not yet collected + * ({@link #pendingTickets}), and open questions not yet answered or lapsed + * ({@link #pendingQuestions}) — and sends at most one combined nudge per tick + * ({@link #injectNudge(String, int, int, int)}). Work that arrives while the lead is busy is + * never lost: it is re-read fresh on every tick until the lead is injectable or its own reminder + * cap ({@link #maxReminders}) is reached — each source spends from its own budget, so one source + * exhausting its cap does not stop nudges about the others (post-CB-590 regression fix; see + * {@link #decide}) — whichever the durable inbox / pending set doesn't already answer via + * {@code STOP}. */ public final class ReplyPushLoop { @@ -57,6 +60,14 @@ public final class ReplyPushLoop { /** Coalesced form, several uncollected tickets for the same lead. */ static final String TICKETS_NUDGE_FORMAT = "%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s"; + /** Singular form, one worker paused mid-turn in bridge_ask (CB-582) — names the answer call directly. */ + static final String QUESTION_NUDGE_FORMAT = + "Worker %s asked a question (ticket %s) — answer it with bridge_send(turnId=\"%s\", " + + "content=...) to resume its turn:\n%s"; + /** Coalesced form, several open questions for the same lead. */ + static final String QUESTIONS_NUDGE_FORMAT = + "%d workers are paused on a question — run bridge_poll(ticket=...) for each, then answer " + + "with bridge_send(turnId=..., content=...): %s"; private final PrimaryRegistry primaryRegistry; private final AgentControl agents; @@ -75,7 +86,14 @@ public final class ReplyPushLoop { private final ConcurrentHashMap pendingReplies = new ConcurrentHashMap<>(); /** Tickets that have gone terminal but not yet been polled, keyed by ticket. */ private final ConcurrentHashMap pendingTickets = new ConcurrentHashMap<>(); - /** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */ + /** + * Open {@code bridge_ask} questions not yet answered or lapsed, keyed by {@code turnId} + * (CB-582). A question's own nudge count is tracked the same per-item way as + * {@link #pendingTickets} (CB-598): a fresh question keeps its source eligible regardless of + * how depleted an older, still-open question's count is. + */ + private final ConcurrentHashMap pendingQuestions = new ConcurrentHashMap<>(); + /** CB-590: leads with an active combined reminder schedule (replies and/or tickets and/or questions). */ private final ConcurrentHashMap activeLeads = new ConcurrentHashMap<>(); public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, @@ -154,6 +172,26 @@ public final class ReplyPushLoop { .collect(Collectors.toUnmodifiableSet()); } + /** + * An open question awaiting the lead's answer: which ticket it belongs to, which worker asked, + * which lead to nudge, the question text, and how many nudges have named it so far (CB-598 — + * tracked per question, not per lead per source). + */ + private record PendingQuestion(String turnId, String ticket, String target, String lead, + String question, int nudgeCount) { + } + + /** Questions still open for {@code lead}, snapshotted fresh for one tick. */ + private List pendingQuestionsFor(String lead) { + return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList(); + } + + /** Question turnIds still open for {@code lead} — a plain snapshot for race comparison. */ + private Set pendingQuestionTurnIdsFor(String lead) { + return pendingQuestionsFor(lead).stream().map(PendingQuestion::turnId) + .collect(Collectors.toUnmodifiableSet()); + } + /** * The reply-source reminder count {@link #decide} should see for {@code lead} on this tick: * the minimum nudge count among the reply targets currently pending for it (CB-598). @@ -189,6 +227,15 @@ public final class ReplyPushLoop { return min == Integer.MAX_VALUE ? 0 : min; } + /** As {@link #minReplyNudgeCountFor}, for the question source (CB-582). */ + private int minQuestionNudgeCountFor(String lead) { + int min = Integer.MAX_VALUE; + for (PendingQuestion q : pendingQuestionsFor(lead)) { + min = Math.min(min, q.nudgeCount()); + } + return min == Integer.MAX_VALUE ? 0 : min; + } + /** * Pure decision function: examine everything pending for {@code lead} — reply targets and * tickets alike — and return what the loop should do. @@ -213,15 +260,29 @@ public final class ReplyPushLoop { * @return the action the caller should take */ Action decide(String lead, int replyReminderCount, int ticketReminderCount) { + return decide(lead, replyReminderCount, ticketReminderCount, minQuestionNudgeCountFor(lead)); + } + + /** + * As {@link #decide(String, int, int)}, with the question source (CB-582) folded in on the + * same footing as replies and tickets: its own eligibility (has open questions AND under its + * own {@link #maxReminders} budget) is enough on its own to {@link Action#INJECT}, exactly like + * the other two. + * + * @param questionReminderCount the lowest nudge count among questions open for this lead + */ + Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount) { boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty(); boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty(); - if (!hasReplyWork && !hasTicketWork) { + boolean hasQuestionWork = !pendingQuestionTurnIdsFor(lead).isEmpty(); + if (!hasReplyWork && !hasTicketWork && !hasQuestionWork) { log.debug("push: nothing pending for lead {}, stopping reminder", lead); return Action.STOP; } boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders; boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders; - if (!replyEligible && !ticketEligible) { + boolean questionEligible = hasQuestionWork && questionReminderCount < maxReminders; + if (!replyEligible && !ticketEligible && !questionEligible) { log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping", maxReminders, lead); countNudge("exhausted"); @@ -300,6 +361,39 @@ public final class ReplyPushLoop { pendingTickets.remove(ticket); } + /** + * Called when an async ticket's worker pauses mid-turn in {@code bridge_ask} (CB-582): the + * question is now visible via {@code bridge_poll} (Phase.ASKING), but the reverse-rendezvous + * window it opened with (~55s default, see {@code BridgeMcp}/{@code BridgedApp}) is far shorter + * than a lead's normal minutes-long poll cadence — exactly the gap this closes. Resolves the + * delegating lead the same way {@link #onTicketTerminal} does and coalesces onto the same + * per-lead schedule (CB-590). + * + * @param ticket the async ticket the question belongs to (for {@code bridge_poll}) + * @param target the worker session that asked + * @param turnId correlation id the lead answers with ({@code bridge_send turnId=...}) + * @param question the question text + */ + public void onQuestionOpened(String ticket, String target, String turnId, String question) { + var lead = primaryRegistry.nudgeTargetFor(target); + if (lead.isEmpty()) { + log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge", + target, turnId); + return; + } + pendingQuestions.put(turnId, new PendingQuestion(turnId, ticket, target, lead.get(), question, 0)); + startOrCoalesce(lead.get()); + } + + /** + * Called when a worker's {@code bridge_ask} resolves — answered or lapsed unanswered — so a + * scheduled tick never nudges about a question the lead already handled. A {@code turnId} that + * was never pending (never nudged, or already closed) is a no-op. + */ + public void questionClosed(String turnId) { + pendingQuestions.remove(turnId); + } + // --- the schedule ---------------------------------------------------------------------------- /** Start a reminder schedule for {@code lead}, or join the one already running. */ @@ -330,17 +424,19 @@ public final class ReplyPushLoop { void tick(String lead) { Set repliesBefore = pendingReplyTargetsFor(lead); Set ticketsBefore = pendingTicketIdsFor(lead); + Set questionsBefore = pendingQuestionTurnIdsFor(lead); int replyReminderCount = minReplyNudgeCountFor(lead); int ticketReminderCount = minTicketNudgeCountFor(lead); - var action = decide(lead, replyReminderCount, ticketReminderCount); + int questionReminderCount = minQuestionNudgeCountFor(lead); + var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount); switch (action) { case INJECT -> { - injectNudge(lead, replyReminderCount, ticketReminderCount); + injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount); scheduleNext(lead); } // Re-check after the configured backoff; the lead may become injectable soon. case WAIT_BUSY -> scheduleNext(lead); - case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore); + case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore); } } @@ -375,9 +471,20 @@ public final class ReplyPushLoop { * side always wins; neither can miss the other, so this never loops on its own account. */ void stopOrRestart(String lead, Set repliesBefore, Set ticketsBefore) { + stopOrRestart(lead, repliesBefore, ticketsBefore, pendingQuestionTurnIdsFor(lead)); + } + + /** + * As {@link #stopOrRestart(String, Set, Set)}, with the question source's (CB-582) own "before" + * snapshot folded into the same race check: a question that raced in during the + * decision-to-release window reclaims the schedule slot exactly like a raced-in reply or ticket. + */ + void stopOrRestart(String lead, Set repliesBefore, Set ticketsBefore, + Set questionsBefore) { activeLeads.remove(lead); boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t)) - || pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t)); + || pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t)) + || pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t)); if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) { log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead); scheduleNext(lead); @@ -387,36 +494,44 @@ public final class ReplyPushLoop { } /** Send one combined nudge covering everything currently pending for {@code lead}. */ - private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount) { - // Re-read rather than threading it down from decide(): a reply can drain, or a ticket be - // collected (or another arrive), between the decision and the injection. + private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount, + int questionReminderCount) { + // Re-read rather than threading it down from decide(): a reply can drain, a ticket be + // collected, or a question be answered (or another arrive), between the decision and the + // injection. Set replyTargets = pendingReplyTargetsFor(lead); List tickets = pendingTicketsFor(lead); - if (replyTargets.isEmpty() && tickets.isEmpty()) { + List questions = pendingQuestionsFor(lead); + if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty()) { log.debug("push: pending work for lead {} drained before the nudge could be sent", lead); return; } - String nudge = formatNudge(replyTargets, tickets); + String nudge = formatNudge(replyTargets, tickets, questions); try { agents.send(lead, nudge); - log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))", + log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; " + + "{} reply target(s), {} ticket(s), {} question(s))", lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, - replyTargets.size(), tickets.size()); + questionReminderCount + 1, maxReminders, + replyTargets.size(), tickets.size(), questions.size()); countNudge("delivered"); } catch (RuntimeException e) { - log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}", - lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString()); + log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}", + lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, + questionReminderCount + 1, maxReminders, e.toString()); } // Bump every item actually named in this nudge, not just whatever the shared source-level // eligibility used to gate (CB-598) — each item's own count is what the next tick's - // minReplyNudgeCountFor / minTicketNudgeCountFor will read. An item already at or over the - // cap keeps riding along in the text (still pending, still named) but its extra bumps here - // are inert: decide() already treats it as ineligible once its count reaches maxReminders. - bumpNudgeCounts(replyTargets, tickets); + // minReplyNudgeCountFor / minTicketNudgeCountFor / minQuestionNudgeCountFor will read. An + // item already at or over the cap keeps riding along in the text (still pending, still + // named) but its extra bumps here are inert: decide() already treats it as ineligible once + // its count reaches maxReminders. + bumpNudgeCounts(replyTargets, tickets, questions); } /** Record that every one of these items was just named in a sent (or attempted) nudge. */ - private void bumpNudgeCounts(Set replyTargets, List tickets) { + private void bumpNudgeCounts(Set replyTargets, List tickets, + List questions) { for (String target : replyTargets) { pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1)); } @@ -424,6 +539,11 @@ public final class ReplyPushLoop { pendingTickets.computeIfPresent(ticket.ticket(), (id, e) -> new PendingTicket(e.ticket(), e.lead(), e.failed(), e.nudgeCount() + 1)); } + for (PendingQuestion question : questions) { + pendingQuestions.computeIfPresent(question.turnId(), (id, e) -> + new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(), + e.nudgeCount() + 1)); + } } /** Schedule the next tick on the scheduler thread pool. */ @@ -435,7 +555,8 @@ public final class ReplyPushLoop { // --- nudge formatting ------------------------------------------------------------------------ /** Render everything pending for one lead as a single nudge line. */ - private static String formatNudge(Set replyTargets, List tickets) { + private static String formatNudge(Set replyTargets, List tickets, + List questions) { List parts = new ArrayList<>(); if (!replyTargets.isEmpty()) { parts.add(formatRepliesNudge(replyTargets)); @@ -443,6 +564,9 @@ public final class ReplyPushLoop { if (!tickets.isEmpty()) { parts.add(formatTicketsNudge(tickets)); } + if (!questions.isEmpty()) { + parts.add(formatQuestionsNudge(questions)); + } return String.join(" | ", parts); } @@ -470,6 +594,18 @@ public final class ReplyPushLoop { return TICKETS_NUDGE_FORMAT.formatted(pending.size(), failedNote, ids); } + /** Render one or several open questions (CB-582). */ + private static String formatQuestionsNudge(List pending) { + if (pending.size() == 1) { + PendingQuestion q = pending.get(0); + return QUESTION_NUDGE_FORMAT.formatted(q.target(), q.ticket(), q.turnId(), q.question()); + } + String ids = pending.stream() + .map(q -> q.ticket() + " (turnId=" + q.turnId() + ")") + .collect(Collectors.joining(", ")); + return QUESTIONS_NUDGE_FORMAT.formatted(pending.size(), ids); + } + // --- lifecycle ----------------------------------------------------------------------------- /** @@ -477,8 +613,8 @@ public final class ReplyPushLoop { * uses this to stand aside: while the push loop is actively nudging a lead, a concurrent * heartbeat injection would start a second competing turn in the same pane — racing loops * multiply turns and context burn. "Active" means a schedule exists in {@link #activeLeads}, - * which now covers both reply-queued (CB-307) and ticket-terminal (CB-588) work (CB-590) — - * bounded by what has been triggered, not by any persistent state. + * which now covers reply-queued (CB-307), ticket-terminal (CB-588), and question-open (CB-582) + * work (CB-590) — bounded by what has been triggered, not by any persistent state. */ public boolean isActive() { return !activeLeads.isEmpty(); @@ -490,6 +626,7 @@ public final class ReplyPushLoop { activeLeads.clear(); pendingReplies.clear(); pendingTickets.clear(); + pendingQuestions.clear(); } /** @see #stop() */ diff --git a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java index 754bb9b..2aa4d7b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -509,10 +509,20 @@ public final class BridgedApp { return; } try { - ctx.status(200).json(Map.of( - "sessionId", id, - "status", messages.status(id).name().toLowerCase(), - "ready", presence.isPresent(id))); + Map body = new LinkedHashMap<>(); + body.put("sessionId", id); + body.put("status", messages.status(id).name().toLowerCase()); + body.put("ready", presence.isPresent(id)); + // CB-582: a worker paused mid-turn in an async bridge_ask is otherwise invisible to a + // status poll — surface the open question and how to answer it, same as bridge_poll's + // Phase.ASKING view. + MessageService.PendingAsk ask = messages.pendingAsk(id); + if (ask != null) { + body.put("question", ask.question()); + body.put("turnId", ask.turnId()); + body.put("ticket", ask.ticket()); + } + ctx.status(200).json(body); } catch (HerdrException e) { herdrError(ctx, e); } @@ -538,6 +548,11 @@ public final class BridgedApp { if (v.detail() != null) { body.put("detail", v.detail()); } + // CB-582: Phase.ASKING carries the question in v.reply() (handled above) and its answer- + // correlation id here — a REST caller polling this ticket otherwise has no way to answer it. + if (v.turnId() != null) { + body.put("turnId", v.turnId()); + } ctx.status(200).json(body); } diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index c584d08..fbbedf4 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -773,6 +773,52 @@ class BridgeMcpTest { assertEquals("blocked", textOf(res)); } + /** + * CB-582: a lead polling {@code bridge_status} on its normal cadence — not {@code bridge_poll} + * — must also see a worker's open async {@code bridge_ask} question, since the reverse-rendezvous + * window it opened with is far shorter than that cadence. + */ + @Test + void statusReportsAnOpenQuestionWhenTheWorkerIsMidAsk() throws Exception { + String ticket = messages.sendAsync(T, "task that asks"); + long deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting(T), "sendAsync should have opened its rendezvous waiter"); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + + MessageService.TaskView asking; + deadline = System.currentTimeMillis() + 3000; + do { + asking = messages.poll(ticket); + Thread.sleep(5); + } while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline); + assertEquals(MessageService.Phase.ASKING, asking.phase()); + + McpSchema.CallToolResult res = BridgeMcp.status(messages, T); + assertNotEquals(Boolean.TRUE, res.isError()); + String out = textOf(res); + assertTrue(out.startsWith("idle"), "the live status must still lead the text: " + out); + assertTrue(out.contains("which config file?"), "the question text must be shown: " + out); + assertTrue(out.contains("turnId=\"" + asking.turnId() + "\""), "the turnId must be shown: " + out); + assertTrue(out.contains("(ticket " + ticket + ")"), "the ticket must be shown: " + out); + + // Clean up the still-open ask so the background thread does not linger past the test. + String turnId = asking.turnId(); + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> messages.answer(turnId, "config.yaml", 5000)); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.resolve(T, "done")); + answer.get(5, TimeUnit.SECONDS); + } + // --- bridge_whoami: the caller's own identity, so an agent never has to guess its role ------- @Test diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 07a0a33..eb2563e 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -714,6 +714,47 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome()); } + // --- CB-582: bridge_status pendingAsk() ------------------------------------------------------ + + @Test + void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception { + assertNull(messages.pendingAsk(T), "no async ticket at all -> no pending ask"); + + String ticket = messages.sendAsync(T, "long task"); + awaitWaiting(); + assertNull(messages.pendingAsk(T), "a plain pending delegation is not a question"); + + injectDelivery(); + assertTrue(rendezvous.resolve(T, "done")); + awaitTicketPhase(ticket, MessageService.Phase.DONE); + assertNull(messages.pendingAsk(T), "a finished ticket carries no open question either"); + } + + @Test + void pendingAskReturnsTheOpenQuestionForAnAsyncTicket() throws Exception { + String ticket = messages.sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); + + MessageService.PendingAsk pending = messages.pendingAsk(T); + assertNotNull(pending, "bridge_status should see the open question"); + assertEquals(ticket, pending.ticket()); + assertEquals("which config file?", pending.question()); + assertEquals(asking.turnId(), pending.turnId()); + + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> messages.answer(asking.turnId(), "config.yaml", 5000)); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + assertNull(messages.pendingAsk(T), "an answered question is no longer pending"); + awaitWaiting(); + assertTrue(rendezvous.resolve(T, "done")); + assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome()); + } + // --- CB-588: async ticket terminal nudges --------------------------------------------------- // // MessageService.reply's rendezvous fast path is exactly what an async ticket always takes @@ -833,6 +874,71 @@ class MessageServiceTest { } } + // --- CB-582: bridge_ask question-open nudges -------------------------------------------------- + + @Test + void anAsyncTicketThatPausesOnAQuestionNudgesTheLeadWithNoPriorPollCall() throws Exception { + try (var wiring = wireWithPushLoop(5, 50)) { + String ticket = wiring.service().sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = CompletableFuture.supplyAsync( + () -> wiring.service().ask(T, "which config file?", 5000)); + MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING); + + awaitNudge(wiring.leadHerdr()); + String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); + assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge); + assertTrue(nudge.contains(asking.turnId()), "the nudge should name the turnId: " + nudge); + assertTrue(nudge.contains("bridge_send(turnId="), + "the nudge should name the exact answer call: " + nudge); + assertTrue(nudge.contains("which config file?"), "the nudge should include the question: " + nudge); + + // Clean up the still-open ask so the background thread does not linger past the test. + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> wiring.service().answer(asking.turnId(), "config.yaml", 5000)); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + awaitWaiting(); + assertTrue(rendezvous.resolve(T, "done")); + answer.get(5, TimeUnit.SECONDS); + } + } + + @Test + void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception { + try (var wiring = wireWithPushLoop(5, 50)) { + String ticket = wiring.service().sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = CompletableFuture.supplyAsync( + () -> wiring.service().ask(T, "which config file?", 5000)); + MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING); + + awaitNudge(wiring.leadHerdr()); + long callsBeforeAnswer = wiring.leadHerdr().calls.stream() + .filter(c -> c.method().equals("agent.prompt")).count(); + + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> wiring.service().answer(asking.turnId(), "config.yaml", 5000)); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + awaitWaiting(); + assertTrue(rendezvous.resolve(T, "done")); + answer.get(5, TimeUnit.SECONDS); + + // Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire — + // none of them may still name the question's turnId, which is closed. + Thread.sleep(300); + boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream() + .filter(c -> c.method().equals("agent.prompt")) + .skip(callsBeforeAnswer) + .anyMatch(c -> c.params().toString().contains(asking.turnId())); + assertFalse(anyNamesClosedQuestion, + "no nudge sent after the answer may still name the now-closed turnId " + asking.turnId()); + } + } + @Test void aFleetWithNoPushLoopConfiguredBehavesExactlyAsToday() throws Exception { // `messages` (the shared field) uses the no-pushLoop constructor — poll() must not throw, diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java index 71557ab..3555fe2 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -657,6 +657,132 @@ class ReplyPushLoopTest { assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh ticket"); } + // --- CB-582: bridge_ask question-open nudges ------------------------------------------------- + + @Test + void onQuestionOpenedWithNoKnownLeadNeverStartsASchedule() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 50); + + loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?"); + + assertFalse(loop.isActive(), "no lead known means nothing to nudge yet"); + Thread.sleep(150); + assertEquals(0, rec.sendCount(), "must not nudge when no lead is known to be waiting"); + } + + @Test + void decideQuestionsWithNothingPendingIsStop() { + agents = agentWithStatus("idle"); + assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0, 0)); + } + + @Test + void decideQuestionsAtCapIsStop() { + agents = agentWithStatus("idle"); + var loop = loop(2, 100_000); + loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?"); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0, 2)); + } + + @Test + void decideQuestionsUnderCapWithInjectableLeadIsInject() { + agents = agentWithStatus("idle"); + var loop = loop(5, 100_000); + loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?"); + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0, 0)); + } + + @Test + void onQuestionOpenedCausesExactlyOneNudgeNamingTheTicketAndTurnId() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + + loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config file?"); + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one question nudge should have been sent"); + assertEquals(1, rec.sendCount()); + String nudge = rec.sentParams().getFirst().getValue().toString(); + assertTrue(nudge.contains("task-1"), "nudge should name the ticket: " + nudge); + assertTrue(nudge.contains("term_worker#1"), "nudge should name the turnId: " + nudge); + assertTrue(nudge.contains("bridge_send(turnId="), "nudge should name the exact answer call: " + nudge); + assertTrue(nudge.contains("which config file?"), "nudge should include the question text: " + nudge); + } + + @Test + void questionClosedPreventsFurtherNudging() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + var loop = loop(1, 100); + + loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?"); + loop.questionClosed("term_worker#1"); // answered/lapsed before the first tick fired + + Thread.sleep(300); // let the scheduled tick run + assertEquals(0, rec.sendCount(), "an already-closed question must never be nudged"); + } + + @Test + void questionNudgesSendUpToCapThenStop() throws Exception { + int cap = 2; + var rec = recordingClient(); + agents = new AgentControl(rec); + rec.sendLatch = new CountDownLatch(cap); + + loop(cap, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?"); + + assertTrue(rec.sendLatch.await(5, TimeUnit.SECONDS), cap + " question nudges should have fired"); + Thread.sleep(300); + assertEquals(cap, rec.sendCount(), "exactly " + cap + " question nudges (cap=" + cap + ")"); + } + + @Test + void questionAndTicketForTheSameLeadCoalesceIntoOneSend() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + var loop = loop(1, 300); // backoff wide enough that both entry points land before the first tick + + loop.onTicketTerminal("task-1", WORKER, false); + loop.onQuestionOpened("task-2", WORKER, "term_worker#1", "which config?"); + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one combined nudge should have been sent"); + Thread.sleep(300); + assertEquals(1, rec.sendCount(), + "a ticket and a question for the same lead must coalesce onto ONE schedule"); + String nudge = rec.sentParams().getFirst().getValue().toString(); + assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge); + assertTrue(nudge.contains("term_worker#1"), "the combined nudge must still mention the question: " + nudge); + } + + @Test + void oneExhaustedQuestionSourceDoesNotBlockANudgeForTheOtherSources() { + // Mirrors oneExhaustedSourceDoesNotBlockANudgeForTheOtherSource for the question source: + // the question source is at its cap (2/2), but the ticket source has never been nudged + // (0/2) — decide() must still INJECT so the ticket is not stranded. + agents = agentWithStatus("idle"); + var loop = loop(2, 100_000); + loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?"); + loop.onTicketTerminal("task-2", WORKER, false); + + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0, 2), + "the question source is exhausted (2/2), but the ticket source has never been " + + "nudged (0/2) — the lead must still be injected so the ticket is not lost"); + } + + @Test + void questionNudgeFormatIsCorrect() { + String single = ReplyPushLoop.QUESTION_NUDGE_FORMAT.formatted( + WORKER, "task-1", "term_worker#1", "which config?"); + assertTrue(single.contains("Worker term_worker")); + assertTrue(single.contains("bridge_send(turnId=\"term_worker#1\"")); + assertTrue(single.contains("which config?")); + + String multi = ReplyPushLoop.QUESTIONS_NUDGE_FORMAT.formatted(2, "task-1 (turnId=t1), task-2 (turnId=t2)"); + assertTrue(multi.contains("2 workers")); + assertTrue(multi.contains("bridge_poll(ticket=...)")); + } + // --- metrics (CB-512) ---------------------------------------------------------------------- @Test diff --git a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java index 1022545..9ae57be 100644 --- a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java @@ -484,6 +484,67 @@ class BridgedAppTest { assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean()); } + /** + * CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the + * ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code bridge_ask} + * question, since the reverse-rendezvous window it opened with is far shorter than that cadence. + */ + @Test + void sessionStatusReportsAnOpenQuestionWhenTheWorkerIsMidAsk() throws Exception { + FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // poller delivers the injection + int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw")); + + HttpResponse accepted = postMessage(port, "{\"content\":\"do it\",\"wait\":false}"); + assertEquals(202, accepted.statusCode()); + String ticket = mapper.readTree(accepted.body()).get("ticket").asText(); + + Thread.sleep(200); // let the background async send open its rendezvous waiter + + var ask = java.util.concurrent.CompletableFuture.supplyAsync(() -> { + try { + return postJson(port, "/sessions/term_a/ask", + "{\"question\":\"which config file?\",\"timeoutMs\":5000}"); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + + JsonNode task; + long deadline = System.currentTimeMillis() + 3000; + do { + task = mapper.readTree(req(port, "GET", "/tasks/" + ticket).body()); + if ("asking".equals(task.path("phase").asText())) break; + //noinspection BusyWait + Thread.sleep(10); + } while (System.currentTimeMillis() < deadline); + assertEquals("asking", task.get("phase").asText()); + String turnId = task.get("turnId").asText(); + + HttpResponse status = req(port, "GET", "/sessions/term_a/status"); + assertEquals(200, status.statusCode()); + JsonNode body = mapper.readTree(status.body()); + assertEquals("idle", body.get("status").asText(), "the live status must still be reported"); + assertEquals("which config file?", body.get("question").asText()); + assertEquals(turnId, body.get("turnId").asText()); + assertEquals(ticket, body.get("ticket").asText()); + + // Answer it — via the same /message route bridge_send uses, keyed by turnId — so the + // background ask thread does not linger past the test. + var answer = java.util.concurrent.CompletableFuture.supplyAsync(() -> { + try { + return postJson(port, "/sessions/term_a/message", + "{\"content\":\"config.yaml\",\"turnId\":\"" + turnId + "\",\"timeoutMs\":4000}"); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + HttpResponse askResponse = ask.get(6, java.util.concurrent.TimeUnit.SECONDS); + assertEquals(200, askResponse.statusCode()); + assertEquals("config.yaml", mapper.readTree(askResponse.body()).get("answer").asText()); + postJson(port, "/sessions/term_a/reply", "{\"content\":\"done\"}"); + answer.get(6, java.util.concurrent.TimeUnit.SECONDS); + } + @Test void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception { FakeHerdr herdr = new FakeHerdr();