Merge worker/t365-3920c5-3: t365: fleet_reply/REST reply distinguish resolved vs queued; rename nudge metric outcome
This commit is contained in:
@@ -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");
|
||||||
|
|||||||
Reference in New Issue
Block a user