CB-582: make a pending bridge_ask question visible on the lead's poll cadence
CI / build (pull_request) Successful in 1m23s
CI / contract (pull_request) Successful in 1m23s

bridge_ask blocks the worker's turn for ~55s by default (BridgedApp.java,
BridgeMcp.java) — a value deliberately kept just under the worker's own MCP
client's ~60s call cap so the daemon can return a clean timeout before the
client severs the call, not a value that can usefully be widened. A lead
following the charter's wait:false + poll cadence is minutes away, so the
window closes long before a poll would ever see the question — and until now
bridge_poll on such a ticket just read as ordinary "pending" progress.

bridge_poll(ticket) already surfaced Phase.ASKING with the question and
turnId (CB-205); this ships the two pieces that were still missing:

- The lead's own pane is now nudged the instant a question opens, reusing
  the CB-588 ReplyPushLoop push mechanism (a third source alongside queued
  replies and terminal tickets) rather than a new path. The nudge is capped
  by the loop's existing maxReminders budget, and stops the moment the
  question is answered or lapses.
- bridge_status(sessionId) and REST GET /sessions/{id}/status now also show
  an open question and how to answer it, via a new
  MessageService.pendingAsk() lookup — covering the case where a lead checks
  status directly rather than the ticket.
- The REST /tasks/{ticket} endpoint was silently missing turnId on an ASKING
  phase (only the MCP layer's formatted text carried it) — fixed as part of
  making the state genuinely visible over both surfaces.

An unanswered question still behaves as today: the worker proceeds and its
reply says the ask went unanswered — not a hard failure.
This commit is contained in:
Dai Ha
2026-08-16 18:35:45 +02:00
parent fdfd4ac491
commit 08968bb1b7
8 changed files with 602 additions and 53 deletions
@@ -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<MessageService.AskResult> 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<MessageService.Reply> 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
@@ -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<MessageService.AskResult> 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<MessageService.Reply> 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<MessageService.AskResult> 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<MessageService.Reply> 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<MessageService.AskResult> 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<MessageService.Reply> 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,
@@ -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
@@ -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<String> 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<String> 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<String> 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();