Merge worker/t365-3920c5-3: t365: fleet_reply/REST reply distinguish resolved vs queued; rename nudge metric outcome
CI / contract (push) Successful in 46s
CI / build (push) Successful in 1m48s

This commit is contained in:
Dai Ha
2026-09-09 07:35:00 +07:00
10 changed files with 159 additions and 45 deletions
@@ -870,6 +870,11 @@ public final class FleetMcp {
* or — when no send is open — queueing the reply in the inbox for later drain (CB-307). * or — when no send is open — queueing the reply in the inbox for later drain (CB-307).
* {@code callerTerminal} is resolved from the connection (never an argument); a {@code null} * {@code callerTerminal} is resolved from the connection (never an argument); a {@code null}
* means the caller is not a known worker (e.g. the primary called it by mistake). * means the caller is not a known worker (e.g. the primary called it by mistake).
*
* <p>fleetd #365: the result text names which of those actually happened
* ({@link MessageService.ReplyOutcome#description()}) instead of the single word "delivered"
* for both — a queued reply is a real success, but it is not the same fact as one that resolved
* a live waiter, and the caller could not previously tell them apart.
*/ */
static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) { static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) {
if (callerTerminal == null) { if (callerTerminal == null) {
@@ -884,8 +889,8 @@ public final class FleetMcp {
if (isBlank(content)) { if (isBlank(content)) {
return error("content is required"); return error("content is required");
} }
messages.reply(callerTerminal, content); MessageService.ReplyOutcome outcome = messages.reply(callerTerminal, content);
return text("delivered"); return text(outcome.description());
} }
/** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */ /** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */
@@ -56,10 +56,13 @@ public final class FleetMetrics {
m.describe(REPLIES, "counter", m.describe(REPLIES, "counter",
"Worker replies by delivery path (rendezvous=resolved an open send, inbox=stranded and held)."); "Worker replies by delivery path (rendezvous=resolved an open send, inbox=stranded and held).");
m.describe(PUSH_NUDGES, "counter", m.describe(PUSH_NUDGES, "counter",
"CB-307 push-loop nudges to the primary (delivered|exhausted)."); "CB-307 push-loop nudges to the primary (sent|exhausted). fleetd #365: \"sent\" means "
+ "the herdr paste-and-submit call succeeded, not that the pane read it — this "
+ "layer has no read-receipt concept.");
m.describe(HEARTBEAT_NUDGES, "counter", m.describe(HEARTBEAT_NUDGES, "counter",
"CB-551 idle-lead heartbeat nudges (delivered|failed|exhausted). Quiet-cap exhaustion " "CB-551 idle-lead heartbeat nudges (sent|failed|exhausted). Quiet-cap exhaustion "
+ "means the lead idled with nothing pending and was told to stand down."); + "means the lead idled with nothing pending and was told to stand down. "
+ "fleetd #365: \"sent\" means the herdr call succeeded, not that the lead read it.");
m.describe(SPAWNS, "counter", m.describe(SPAWNS, "counter",
"Worker spawn attempts by peer kind and outcome (ready|timeout|guard_rejected)."); "Worker spawn attempts by peer kind and outcome (ready|timeout|guard_rejected).");
m.describe(HERDR_CALLS, "counter", m.describe(HERDR_CALLS, "counter",
@@ -96,7 +96,14 @@ public final class LeadHeartbeatLoop {
this.metrics = metrics; this.metrics = metrics;
} }
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */ /**
* Count one nudge outcome when a registry is wired; a no-op in unit tests.
*
* <p>fleetd #365: the {@code "sent"} outcome (renamed from {@code "delivered"}) records only
* that {@link #injectNudge} — a one-way herdr {@code agent.prompt} paste-and-submit — returned
* without throwing, not that the lead's pane actually read or acted on the text. This layer has
* no read-receipt concept, so "sent" is the honest word for what this call can ever establish.
*/
private void countNudge(String outcome) { private void countNudge(String outcome) {
if (metrics != null) { if (metrics != null) {
metrics.inc(FleetMetrics.HEARTBEAT_NUDGES, "outcome", outcome); metrics.inc(FleetMetrics.HEARTBEAT_NUDGES, "outcome", outcome);
@@ -245,7 +252,7 @@ public final class LeadHeartbeatLoop {
agents.send(leadTerminal, fleet.nudgeText()); agents.send(leadTerminal, fleet.nudgeText());
log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})", log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})",
leadTerminal, quietCount); leadTerminal, quietCount);
countNudge("delivered"); countNudge("sent");
} catch (RuntimeException e) { } catch (RuntimeException e) {
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString()); log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
countNudge("failed"); countNudge("failed");
@@ -138,6 +138,53 @@ public final class MessageService {
public record AskResult(AskOutcome outcome, String answer) { public record AskResult(AskOutcome outcome, String answer) {
} }
/**
* How a worker's {@code fleet_reply} ({@link #reply(String, String)}) actually landed
* (fleetd #365) — the two doors that expose it, {@code fleet_reply} and {@code POST
* /sessions/{id}/reply}, both used to report the single word "delivered" whichever of these
* happened, so a caller could not tell an active handoff from a reply merely held for later
* drain. Both are successes; they are not the same fact.
*/
public enum ReplyOutcome {
/** Resolved a {@code fleet_send}/{@code fleet_ask} that was actively waiting on this reply. */
RESOLVED_SEND("resolved_send", true,
"delivered — resolved the fleet_send that was waiting for it"),
/**
* No live waiter was open, but the reply completed a parked async ticket directly
* ({@link #askAnsweredAsyncTasks}) — a {@code fleet_poll} caller sees it immediately.
*/
RESOLVED_ASYNC_TICKET("resolved_async_ticket", true,
"delivered — resolved a pending async ticket (visible to fleet_poll)"),
/** Nothing was waiting; the reply was queued in the inbox for a later drain (CB-307). */
QUEUED("queued", false,
"queued — no send or ticket was waiting; held in the inbox for a later drain");
private final String wireName;
private final boolean delivered;
private final String description;
ReplyOutcome(String wireName, boolean delivered, String description) {
this.wireName = wireName;
this.delivered = delivered;
this.description = description;
}
/** Stable machine-readable name for a JSON/metrics label (REST's {@code outcome} field). */
public String wireName() {
return wireName;
}
/** Whether something was actively waiting and received this reply right now. */
public boolean delivered() {
return delivered;
}
/** Shared human-readable text — the one place both {@code fleet_reply} and REST word this. */
public String description() {
return description;
}
}
/** Lifecycle phase of an async delegation ticket. */ /** Lifecycle phase of an async delegation ticket. */
public enum Phase { public enum Phase {
/** Delegated and in flight — queued for the worker or being worked. */ /** Delegated and in flight — queued for the worker or being worked. */
@@ -439,16 +486,15 @@ public final class MessageService {
* @throws IllegalArgumentException if {@code content} is {@code null} or blank — the caller must * @throws IllegalArgumentException if {@code content} is {@code null} or blank — the caller must
* report this as a client error (REST: 400 {@code bad_request}) rather than resolve * report this as a client error (REST: 400 {@code bad_request}) rather than resolve
* anything * anything
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or * @return which of the three ways (fleetd #365) the reply actually landed — never {@code null}
* was queued
*/ */
public boolean reply(String session, String content) { public ReplyOutcome reply(String session, String content) {
if (content == null || content.isBlank()) { if (content == null || content.isBlank()) {
throw new IllegalArgumentException("content is required"); throw new IllegalArgumentException("content is required");
} }
if (rendezvous.resolve(session, content)) { if (rendezvous.resolve(session, content)) {
count(FleetMetrics.REPLIES, "path", "rendezvous"); count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path return ReplyOutcome.RESOLVED_SEND; // a live send took it — unchanged fast path
} }
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming // #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on: // a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
@@ -480,7 +526,7 @@ public final class MessageService {
asyncTasksByTurn.remove(turnId, orphan); asyncTasksByTurn.remove(turnId, orphan);
} }
count(FleetMetrics.REPLIES, "path", "async-recovered"); count(FleetMetrics.REPLIES, "path", "async-recovered");
return true; // the ticket itself took it — no inbox stranding at all return ReplyOutcome.RESOLVED_ASYNC_TICKET; // the ticket itself took it — no inbox stranding
} }
} else if (candidates.size() > 1) { } else if (candidates.size() > 1) {
List<String> tickets = candidates.stream().map(t -> t.ticket).toList(); List<String> tickets = candidates.stream().map(t -> t.ticket).toList();
@@ -498,7 +544,7 @@ public final class MessageService {
if (pushLoop != null) { if (pushLoop != null) {
pushLoop.onReplyQueued(session); pushLoop.onReplyQueued(session);
} }
return true; // held, not lost return ReplyOutcome.QUEUED; // held, not lost — but not delivered either
} }
/** /**
@@ -132,7 +132,15 @@ public final class ReplyPushLoop {
this.metrics = metrics; this.metrics = metrics;
} }
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */ /**
* Count one nudge outcome when a registry is wired; a no-op in unit tests.
*
* <p>fleetd #365: the {@code "sent"} outcome (renamed from {@code "delivered"}) records only
* that {@code agents.send} — a one-way herdr {@code agent.prompt} paste-and-submit — returned
* without throwing. Nothing in this loop, or anywhere downstream of it, confirms the pane
* actually read or acted on the text; there is no read-receipt concept at this layer. "Sent"
* says exactly that; "delivered" claimed more than this call can ever establish.
*/
private void countNudge(String outcome) { private void countNudge(String outcome) {
if (metrics != null) { if (metrics != null) {
metrics.inc(FleetMetrics.PUSH_NUDGES, "outcome", outcome); metrics.inc(FleetMetrics.PUSH_NUDGES, "outcome", outcome);
@@ -730,7 +738,7 @@ public final class ReplyPushLoop {
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders, questionReminderCount + 1, maxReminders,
replyTargets.size(), tickets.size(), questions.size()); replyTargets.size(), tickets.size(), questions.size());
countNudge("delivered"); countNudge("sent");
for (PendingIncident incident : incidents) { for (PendingIncident incident : incidents) {
if (pendingIncidents.remove(incident.key(), incident)) { if (pendingIncidents.remove(incident.key(), incident)) {
deliveredIncidents.add(incident.key()); deliveredIncidents.add(incident.key());
@@ -696,6 +696,10 @@ public final class FleetApp {
/** /**
* The worker's structured reply ({@code fleet_reply}) — resolves the blocking send awaiting * The worker's structured reply ({@code fleet_reply}) — resolves the blocking send awaiting
* on this session, or queues the reply in the inbox when no send is open (CB-307). * on this session, or queues the reply in the inbox when no send is open (CB-307).
*
* <p>fleetd #365: the response body's {@code delivered} field used to be unconditionally
* {@code true} for either case; it now reports whether a send/ticket was actually resolved,
* with {@code outcome} naming which (see {@link MessageService.ReplyOutcome}).
*/ */
private void replyMessage(Context ctx) { private void replyMessage(Context ctx) {
String id = ctx.pathParam("id"); String id = ctx.pathParam("id");
@@ -719,13 +723,19 @@ public final class FleetApp {
// a WRONG value instead of failing loudly. The check lives in MessageService.reply so both // a WRONG value instead of failing loudly. The check lives in MessageService.reply so both
// this door and FleetMcp.reply inherit the same rule; this catch only translates it into the // this door and FleetMcp.reply inherit the same rule; this catch only translates it into the
// {error, detail} envelope this file uses everywhere else. // {error, detail} envelope this file uses everywhere else.
MessageService.ReplyOutcome outcome;
try { try {
messages.reply(id, content); outcome = messages.reply(id, content);
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", e.getMessage())); ctx.status(400).json(Map.of("error", "bad_request", "detail", e.getMessage()));
return; return;
} }
ctx.status(200).json(Map.of("sessionId", id, "delivered", true)); // fleetd #365: "delivered": true used to be unconditional here, whether the reply resolved
// a waiting send or was merely queued in the inbox for a later drain — the same gap
// FleetMcp.reply had over MCP. `delivered` now reflects which actually happened, and
// `outcome` names the specific case (see MessageService.ReplyOutcome).
ctx.status(200).json(Map.of("sessionId", id, "delivered", outcome.delivered(),
"outcome", outcome.wireName()));
} }
/** /**
@@ -108,8 +108,10 @@ class FleetMcpTest {
} }
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter"); assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
// fleetd #365: a resolved live send must read distinctly from a merely-queued reply —
// see replyWithNoPendingSendIsQueuedNotError below for the other case.
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "LGTM"); McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "LGTM");
assertEquals("delivered", textOf(reply)); assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS); McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS);
assertNotEquals(Boolean.TRUE, res.isError()); assertNotEquals(Boolean.TRUE, res.isError());
@@ -135,7 +137,7 @@ class FleetMcpTest {
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter"); assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "async LGTM"); McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "async LGTM");
assertEquals("delivered", textOf(reply)); assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
// Poll until the async send completes and reports the reply. // Poll until the async send completes and reports the reply.
McpSchema.CallToolResult polled = FleetMcp.poll(messages, ticket, null); McpSchema.CallToolResult polled = FleetMcp.poll(messages, ticket, null);
@@ -328,9 +330,10 @@ class FleetMcpTest {
@Test @Test
void replyWithNoPendingSendIsQueuedNotError() { void replyWithNoPendingSendIsQueuedNotError() {
// CB-307: a reply with no open send is now queued in the inbox, not an error. // CB-307: a reply with no open send is now queued in the inbox, not an error.
// fleetd #365: it must also no longer claim "delivered" — nothing was waiting for it.
McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", "orphan"); McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", "orphan");
assertNotEquals(Boolean.TRUE, res.isError(), "a queued reply is not an error"); assertNotEquals(Boolean.TRUE, res.isError(), "a queued reply is not an error");
assertEquals("delivered", textOf(res)); assertEquals(MessageService.ReplyOutcome.QUEUED.description(), textOf(res));
// The reply is drainable by target. // The reply is drainable by target.
var drained = messages.drainReplies("term_a"); var drained = messages.drainReplies("term_a");
@@ -413,7 +416,7 @@ class FleetMcpTest {
} }
assertTrue(rendezvous.isWaiting("term_a"), "the answer should have reopened a waiter"); assertTrue(rendezvous.isWaiting("term_a"), "the answer should have reopened a waiter");
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "done"); McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "done");
assertEquals("delivered", textOf(reply)); assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS))); assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
} }
@@ -479,7 +479,8 @@ class MessageServiceTest {
// The worker resumes on its own (per the ask() contract) and eventually sends its real // The worker resumes on its own (per the ask() contract) and eventually sends its real
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING. // fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted"); assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET, messages.reply(T, "real result"),
"the worker's real reply must still be accepted, resolving the parked async ticket");
} finally { } finally {
messages.setAskTimeoutRaceHookForTest(null); messages.setAskTimeoutRaceHookForTest(null);
} }
@@ -767,7 +768,9 @@ class MessageServiceTest {
@Test @Test
void replyQueuesInInboxWhenNoSendIsOpen() { void replyQueuesInInboxWhenNoSendIsOpen() {
// No send is open for this session — reply should queue in the inbox. // No send is open for this session — reply should queue in the inbox.
assertTrue(messages.reply(T, "queued-text"), "reply should succeed (queued)"); // fleetd #365: this is the case that must read as QUEUED, not "delivered".
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "queued-text"),
"reply should succeed but only as queued — nothing was waiting for it");
var drained = messages.drainReplies(T); var drained = messages.drainReplies(T);
assertEquals(1, drained.size()); assertEquals(1, drained.size());
@@ -780,7 +783,9 @@ class MessageServiceTest {
awaitUninterruptibly(T); awaitUninterruptibly(T);
// An explicit reply resolves the open send. // An explicit reply resolves the open send.
assertTrue(messages.reply(T, "send-resolved"), "reply should succeed (resolved live send)"); // fleetd #365: this is the other case — RESOLVED_SEND, distinct from QUEUED above.
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, messages.reply(T, "send-resolved"),
"reply should succeed by resolving the live waiting send");
// The inbox should be empty — the reply went to the send, not the inbox. // The inbox should be empty — the reply went to the send, not the inbox.
assertTrue(messages.drainReplies(T).isEmpty(), "no reply in the inbox"); assertTrue(messages.drainReplies(T).isEmpty(), "no reply in the inbox");
@@ -1094,8 +1099,11 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(), assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait gives up before the worker finishes resuming"); "the primary's own bounded wait gives up before the worker finishes resuming");
// The worker keeps working past that window and only now calls fleet_reply. // The worker keeps working past that window and only now calls fleet_reply. The forward
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42")); // waiter answer() opened already timed out, so this resolves via the parked async ticket,
// not a live send (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE); MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(), assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1119,7 +1127,10 @@ class MessageServiceTest {
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome()); assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42")); // No live waiter (answer()'s own forward wait already timed out) — resolves the parked
// async ticket instead (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
// fleet_stop tears the worker's session down right after the reply landed — this must never // fleet_stop tears the worker's session down right after the reply landed — this must never
// report the misleading "the worker session was released before it replied": a reply is // report the misleading "the worker session was released before it replied": a reply is
@@ -1170,7 +1181,9 @@ class MessageServiceTest {
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42")); // A live waiter is open (the forward wait above) — this resolves it directly (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND,
messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(), assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the lead's own answer() call must not throw because ask()'s timeout cleanup raced it"); "the lead's own answer() call must not throw because ask()'s timeout cleanup raced it");
@@ -1221,8 +1234,9 @@ class MessageServiceTest {
// The worker keeps working past the timeout and only now calls fleet_reply — with no live // The worker keeps working past the timeout and only now calls fleet_reply — with no live
// rendezvous waiter open (ask()'s timeout already closed it) and no new send() having // rendezvous waiter open (ask()'s timeout already closed it) and no new send() having
// reopened one for this target. // reopened one for this target. So it resolves the parked async ticket (fleetd #365).
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42")); assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE); MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(), assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1279,7 +1293,10 @@ class MessageServiceTest {
// real reply — reproduce that interleaving directly instead of trying to win a real race. // real reply — reproduce that interleaving directly instead of trying to win a real race.
messages.forgetTurnForTest(turnId); messages.forgetTurnForTest(turnId);
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42")); // answer() is still waiting on its own forward waiter for the resumed turn — a live send —
// so this resolves it directly, not the async ticket (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND,
messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(), assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the primary's own answer() call must still see the worker's real reply"); "the primary's own answer() call must still see the worker's real reply");
@@ -1321,7 +1338,9 @@ class MessageServiceTest {
messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId)); messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
try { try {
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42")); // No live waiter — resolves the parked async ticket (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE); MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(), assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1353,7 +1372,8 @@ class MessageServiceTest {
injectDelivery(); injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome()); assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome());
assertTrue(messages.reply(T, "which task does this answer?")); // Ambiguous — two candidates, so it must fall back to the inbox rather than guess (fleetd #365).
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "which task does this answer?"));
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(), assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(),
"an ambiguous reply must not guess ticket1"); "an ambiguous reply must not guess ticket1");
@@ -2032,7 +2052,7 @@ class MessageServiceTest {
CompletableFuture<MessageService.Reply> send = sendAsync(); CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting(); awaitWaiting();
assertTrue(messages.reply(T, "resolved-live")); assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, messages.reply(T, "resolved-live"));
assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded"); assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded");
MessageService.Reply r = send.get(5, TimeUnit.SECONDS); MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
@@ -2042,14 +2062,14 @@ class MessageServiceTest {
@Test @Test
void hasStrandedReplyIsTrueWhenNoSendWasWaiting() { void hasStrandedReplyIsTrueWhenNoSendWasWaiting() {
// No send is open for T — the reply queues into the inbox and is recorded as stranded. // No send is open for T — the reply queues into the inbox and is recorded as stranded.
assertTrue(messages.reply(T, "nobody was waiting")); assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "nobody was waiting"));
assertTrue(messages.hasStrandedReply(T), assertTrue(messages.hasStrandedReply(T),
"a reply with no open send strands, even though it is safely queued in the inbox"); "a reply with no open send strands, even though it is safely queued in the inbox");
} }
@Test @Test
void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception { void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
assertTrue(messages.reply(T, "stray")); assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "stray"));
assertTrue(messages.hasStrandedReply(T)); assertTrue(messages.hasStrandedReply(T));
// The next accepted delivery for T clears the stale stranding fact — the one case the // The next accepted delivery for T clears the stale stranding fact — the one case the
@@ -2067,7 +2087,7 @@ class MessageServiceTest {
@Test @Test
void hasStrandedReplyClearsOnAbandon() { void hasStrandedReplyClearsOnAbandon() {
assertTrue(messages.reply(T, "stray")); assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "stray"));
assertTrue(messages.hasStrandedReply(T)); assertTrue(messages.hasStrandedReply(T));
messages.abandon(T, "session released"); messages.abandon(T, "session released");
@@ -984,7 +984,7 @@ class ReplyPushLoopTest {
// --- metrics (CB-512) ---------------------------------------------------------------------- // --- metrics (CB-512) ----------------------------------------------------------------------
@Test @Test
void successfulNudgeIncrementsDelivered() throws Exception { void successfulNudgeIncrementsSent() throws Exception {
var rec = recordingClient(); var rec = recordingClient();
agents = new AgentControl(rec); agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello"); inbox.publish(WORKER, "m1", "hello");
@@ -994,11 +994,13 @@ class ReplyPushLoopTest {
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"one nudge (1 agent.prompt call) should have been sent"); "one nudge (1 agent.prompt call) should have been sent");
// The delivered count is bumped on the scheduler thread right after the send that releases // The sent count is bumped on the scheduler thread right after the send that releases
// the latch — settle briefly so the counter is published before we read it. // the latch — settle briefly so the counter is published before we read it.
Thread.sleep(200); Thread.sleep(200);
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"), // fleetd #365: "sent", not "delivered" — this only proves the herdr call succeeded, not
"a successfully sent nudge must count as delivered"); // that the primary's pane read it.
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"),
"a successfully sent nudge must count as sent");
} }
@Test @Test
@@ -1013,11 +1015,11 @@ class ReplyPushLoopTest {
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "exhausted"), assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted"); "hitting the reminder cap must count as exhausted");
assertEquals(0, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered")); assertEquals(0, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"));
} }
@Test @Test
void successfulTicketNudgeIncrementsDelivered() throws Exception { void successfulTicketNudgeIncrementsSent() throws Exception {
var rec = recordingClient(); var rec = recordingClient();
agents = new AgentControl(rec); agents = new AgentControl(rec);
Metrics metrics = new Metrics(); Metrics metrics = new Metrics();
@@ -1026,8 +1028,8 @@ class ReplyPushLoopTest {
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one ticket nudge should have been sent"); assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one ticket nudge should have been sent");
Thread.sleep(200); Thread.sleep(200);
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"), assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"),
"a successfully sent ticket nudge must count as delivered, same metric as CB-307"); "a successfully sent ticket nudge must count as sent, same metric as CB-307");
} }
@Test @Test
@@ -476,6 +476,11 @@ class FleetAppTest {
Thread.sleep(200); Thread.sleep(200);
HttpResponse<String> reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}"); HttpResponse<String> reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}");
assertEquals(200, reply.statusCode()); assertEquals(200, reply.statusCode());
// fleetd #365: "delivered" used to be unconditionally true; a live send was actually waiting
// here, so this is the case where it must genuinely read true, with outcome naming why.
JsonNode replyBody = mapper.readTree(reply.body());
assertEquals(true, replyBody.get("delivered").asBoolean());
assertEquals("resolved_send", replyBody.get("outcome").asText());
HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS); HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS);
assertEquals(200, res.statusCode()); assertEquals(200, res.statusCode());
@@ -527,6 +532,11 @@ class FleetAppTest {
int port = startHealthy(); int port = startHealthy();
HttpResponse<String> res = postJson(port, "/sessions/term_a/reply", "{\"content\":\"orphan\"}"); HttpResponse<String> res = postJson(port, "/sessions/term_a/reply", "{\"content\":\"orphan\"}");
assertEquals(200, res.statusCode()); assertEquals(200, res.statusCode());
// fleetd #365: nothing was waiting, so "delivered" must now read false, not the old
// unconditional true — outcome names this as queued.
JsonNode resBody = mapper.readTree(res.body());
assertEquals(false, resBody.get("delivered").asBoolean());
assertEquals("queued", resBody.get("outcome").asText());
// The queued reply is drainable. // The queued reply is drainable.
HttpResponse<String> drain = req(port, "GET", "/sessions/term_a/replies"); HttpResponse<String> drain = req(port, "GET", "/sessions/term_a/replies");