Compare commits

...

8 Commits

Author SHA1 Message Date
Dai Ha bbbb4c1eb3 fleetd #302: require content in MessageService.reply so a REST reply with no content cannot silently resolve a waiter
CI / build (pull_request) Successful in 1m44s
CI / contract (pull_request) Successful in 3m17s
2026-09-04 12:31:09 +07:00
Dai Ha 38c248e617 Merge #298: a released reply is requeued, not dropped
CI / contract (push) Successful in 1m27s
CI / build (push) Successful in 1m46s
release() cancelled the consumer and dropped its local record of deliveries the
broker still held as outstanding. Cancelling a consumer does not requeue them:
they stay unacked on the still-open channel until it or the connection closes.
So a held-but-undrained reply became permanently unreachable — a worker's
report lost with no error and no log line. It now nacks with requeue, after
cancelling, so a later own() can still receive it.
2026-09-04 12:21:30 +07:00
Dai Ha 9debc0de27 #297: render both profiles doors from one body builder, not two copies
CI / contract (push) Successful in 1m23s
CI / build (push) Successful in 1m31s
The merged change shared the QuarantineSource and OutageSource instances
between fleet_profiles and GET /profiles, so the two doors read identical
facts. It then rendered those facts through a character-for-character copy of
the loop, in a different file. Shared inputs do not make duplicated computation
safe: a later edit to the row shape lands on one door and not the other, and
the two disagree about a live outage. That is what #284 was.

The ticket caused this. It said 'read from the same shared instances' and 'do
not change FleetMcp', and together those made copying the loop the only legal
move. Extracting FleetMcp.profilesView and calling it from both is what the
ticket should have asked for.
2026-09-04 12:19:02 +07:00
Dai Ha f34361b263 Merge #297: REST reports what MCP reports
GET /agents and GET /members now map a herdr transport failure into the same
{error, detail} envelope every other handler in FleetApp uses, instead of
letting it escape to Javalin's default handling. GET /profiles now reports the
quarantined and coolingOff states, read from the same shared sources FleetMcp
reads. REST is the door a lead falls back to when its MCP mount drops, so it
was weakest exactly when it was load-bearing.
2026-09-04 12:16:00 +07:00
Dai Ha be5ba22c75 #298: AmqpReplyInbox.release() requeues held deliveries instead of dropping them
CI / contract (pull_request) Successful in 1m12s
CI / build (pull_request) Successful in 2m3s
release(target) used to cancel the target's consumer and clear held's local
record for it. Cancelling a consumer does not requeue the broker's in-flight
deliveries — they stay unacked on the still-open channel until a real
connection drop. So a held-but-undrained reply became permanently
unreachable: never acked, never nacked, never requeued, invisible to peek.

Fix: cancel the consumer first (so it can no longer receive redeliveries),
then nack-with-requeue every held delivery for that target before dropping
the local record. Nacking before the cancel was tried first but a real
broker demonstrated a race: the still-active consumer immediately received
the requeued message back, racing held.remove and leaving peek non-empty.
Cancelling first avoids that. A failed requeue is logged at WARN and does
not abort release(), matching the best-effort teardown style #293 settled
for HerdrPeerLauncher.stop().

Extends AmqpReplyInboxContractTest.releaseCancelsConsumer... to assert the
held delivery is recoverable via a later own(), not just absent from peek.
2026-09-04 12:14:00 +07:00
Dai Ha e97502d550 drainAll: record that the timeout is a whole-drain budget, not a per-session grace
CI / build (push) Successful in 1m27s
CI / contract (push) Successful in 1m48s
The javadoc said 'for each session that is BUSY, poll up to timeoutNanos',
which reads as a per-session grace period. The deadline is taken once, before
the loop, so the first BUSY session can spend all of it. That is deliberate and
is the safer of the two designs: the drain is one phase of a shutdown sequence
that must finish inside launchd's exit window, and a per-session grace would
overrun it and get the daemon SIGKILLed part-way through, leaving the sessions
not yet reached with no clean release, no preserved-worktree log and no
snapshot. Found by a read-only hunt that read the code correctly and drew the
opposite conclusion from the wording.
2026-09-04 12:09:44 +07:00
Dai Ha aef14ff46e Merge #296: a failed spawn no longer leaks the pane it opened
Two exits created a pane and left it running. The readiness gate propagated an
unrelated herdr error without teardown, and spawnAsPane never closed the pane it
split when the peer failed to start. Neither could be cleaned up by the caller:
SessionManager.acquire never learns the pane id, because spawn throws before it
returns one. It removed the worktree anyway, so the leak was a live backend with
a deleted cwd, invisible to fleet_list and holding a seat nothing decremented.
2026-09-04 12:08:23 +07:00
Dai Ha cba516bda4 fleetd #296: close panes on failed spawn
CI / contract (pull_request) Successful in 1m38s
CI / build (pull_request) Successful in 2m7s
2026-09-04 12:05:26 +07:00
10 changed files with 253 additions and 52 deletions
@@ -1015,6 +1015,25 @@ public final class FleetMcp {
* both maps at once when it is both exhaustion-quarantined AND cooling off.
*/
static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
return text(json(profilesView(workers, quarantine, outage)));
}
/**
* The body both front doors answer {@code profiles} with: the configured profile names, the
* default, and the two independent outage states — {@code quarantined} (the backend reported it
* out of capacity) and {@code coolingOff} (the credential threw repeated non-exhaustion backend
* errors). Each map is present only when at least one profile is in that state, and a profile
* can appear in both at once, because the two checks are separate.
*
* <p>fleetd #297: extracted so {@code fleet_profiles} and {@code GET /profiles} render from ONE
* body builder rather than two copies. Passing both doors the same {@link QuarantineSource} and
* {@link OutageSource} instances is necessary but not sufficient: with the loop written out
* twice, a later edit to the row shape — a renamed key, an added field — lands on one door and
* not the other, and the two then disagree about a live outage. That is exactly what fleetd
* #284 was, where one rule computed in two places was widened in only one and a single response
* contradicted itself. Shared inputs do not make duplicated computation safe.
*/
public static Map<String, Object> profilesView(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
Map<String, Object> result = new LinkedHashMap<>();
result.put("profiles", workers.profiles());
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
@@ -1046,7 +1065,7 @@ public final class FleetMcp {
if (!coolingOff.isEmpty()) {
result.put("coolingOff", coolingOff);
}
return text(json(result));
return result;
}
/**
@@ -697,7 +697,20 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
if (paneId == null) {
throw new IllegalStateException("pane.split returned no pane — cannot start a peer");
}
Agent peer = startUniquelyNamed(cfg, argv, paneId).agent();
Agent peer;
try {
peer = startUniquelyNamed(cfg, argv, paneId).agent();
} catch (RuntimeException e) {
// The peer never started — don't leave the pane we just created orphaned.
// Best-effort cleanup; never let it mask the real spawn failure.
try {
stop(paneId);
} catch (RuntimeException cleanup) {
log.warn("failed to close orphaned pane {} after spawn error: {}",
paneId, cleanup.getMessage());
}
throw e;
}
log.info("{} started pane={} terminal={}", namePrefix, peer.paneId(), peer.terminalId());
return peer;
}
@@ -996,7 +1009,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* that gap: it stops waiting immediately (never burns the rest of the timeout), runs the same
* teardown the timeout path below runs, and throws with a message that says the backend exited
* rather than that the pane was slow. Any other {@link HerdrException} still propagates
* unchanged — this gate does not know how to recover from it.
* unchanged — this gate does not interpret or recover from it, but it still closes the pane
* it opened before handing the exception to its caller.
*/
private void waitUntilInjectableOrThrow(String paneId) {
long start = nowMillis.getAsLong();
@@ -1010,7 +1024,15 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
if (isAlreadyGone(e)) {
failFastOnGoneBackend(paneId, e, nowMillis.getAsLong() - start);
}
throw e; // any other herdr failure is not ours to interpret — let it propagate
// This gate must not interpret an unrelated herdr error, but the caller does not
// receive paneId when spawn throws. Close the pane here before propagating e unchanged.
try {
stop(paneId);
} catch (RuntimeException cleanup) {
log.warn("failed to close orphaned pane {} after readiness-gate error: {}",
paneId, cleanup.getMessage());
}
throw e;
}
lastStatus = sample.status();
if (lastStatus.injectable() || refinedInjectable(paneId, sample)) {
@@ -208,18 +208,72 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
/**
* Release ownership of {@code target}: cancel its consumer, then nack-with-requeue every
* delivery still held for it instead of just dropping the local record.
*
* <p><strong>Cancelling a consumer does not requeue its in-flight deliveries.</strong> In AMQP,
* a delivery that was pushed to a consumer stays unacked, attached to the still-open
* {@link #channel}, until that channel or the connection closes — {@code basicCancel} alone does
* neither. So before this method existed with a requeue step, it dropped {@link #held}'s entries
* for {@code target} while the broker still considered them outstanding: never acked, never
* nacked, never requeued, and no longer reachable by {@link #peek} — permanently invisible. This
* is unlike {@link #handleRecovery} and {@link #close()}, whose bare {@code held.clear()} is
* correct because each has already made the broker requeue (a real connection drop, or
* {@code channel.close()} respectively) before clearing local state.
*
* <p><strong>Order: cancel first, then nack.</strong> A delivery tag stays valid for
* {@code basicNack} on this channel regardless of whether its consumer is still attached — only
* a channel/connection close invalidates it — so cancelling {@code target}'s consumer first does
* not risk the tags. Doing it the other way round does: nacking a delivery with {@code requeue}
* while its consumer is still active hands the message straight back to that <em>same</em>
* consumer the instant a prefetch slot frees up (confirmed against a real broker — see
* {@code AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery}),
* which races this method's own {@code held.remove(target)}: the redelivery can land after the
* clear and leave a stale entry behind, so {@link #peek} is no longer reliably empty right after
* {@link #release}. Cancelling first closes that consumer, so the requeued message goes back to
* the queue for whichever consumer picks it up next (a later {@link #own}), not this one.
*
* <p><strong>Failure of the requeue is best-effort, not fatal.</strong> {@link #release} runs
* during teardown ({@code Fleetd} calls it right after {@code MessageService.abandon}), and a
* throw here would abort cleanups the caller depends on — the same argument fleetd #293 settled
* for {@code HerdrPeerLauncher.stop()}'s tab-close step. So a failed {@code basicNack} is logged
* at WARN, naming the target and delivery tag that leaked, and release proceeds; the delivery
* stays unacked on the broker rather than being silently dropped, so it is still recoverable by a
* later connection drop even though this release did not manage to requeue it immediately. A
* failed {@code basicCancel} still throws, unchanged from before this fix — that failure means
* the consumer may still be attached, so best-effort requeue is not attempted underneath it.
*/
@Override
public void release(String target) {
synchronized (channelLock) {
String tag = consumerTags.remove(target);
held.remove(target); // stale delivery tags must not survive release
if (tag == null) {
return;
if (tag != null) {
try {
channel.basicCancel(tag);
} catch (IOException e) {
throw new IllegalStateException("cannot cancel consumer for " + target, e);
}
}
try {
channel.basicCancel(tag);
} catch (IOException e) {
throw new IllegalStateException("cannot cancel consumer for " + target, e);
var perTarget = held.remove(target);
if (perTarget != null) {
synchronized (perTarget) {
for (Held h : perTarget.values()) {
try {
channel.basicNack(h.deliveryTag(), false, true); // requeue, don't drop
} catch (IOException | RuntimeException e) {
// Caught broadly (not just IOException) for the same reason #293 catches
// RuntimeException in HerdrPeerLauncher.stop(): best-effort teardown must
// not be guarded only against the expected failure and bare against any
// other. The message stays unacked on the broker either way — not lost,
// just not proactively requeued — until a connection drop frees it.
log.warn("release({}): could not requeue held delivery (msgId={}, tag={})"
+ " back to the broker — it stays unacked until a connection"
+ " drop frees it: {}",
target, h.message().msgId(), h.deliveryTag(), e.getMessage());
}
}
}
}
}
}
@@ -408,10 +408,28 @@ public final class MessageService {
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
* deterministically instead.
*
* <p><strong>{@code content} is required (fleetd #302).</strong> Both doors that reach this
* method must reject a missing/blank reply the same way, so the check lives here rather than in
* either caller: {@code FleetMcp.reply} already refuses a {@code null} content before it ever
* calls this method (its own required-arg guard), and no test or production call site anywhere
* in the codebase relies on replying with empty content — confirmed by searching every call site
* of this method before adding the check, not assumed. Without this guard, a REST {@code
* POST /sessions/{id}/reply} whose body omits {@code content} (or a client library that maps a
* missing field to {@code ""}) used to reach {@link Rendezvous#resolve} with an empty string,
* silently completing the lead's blocking wait with nothing — indistinguishable from a worker
* that genuinely replied with nothing, which is worse than a loud failure because it destroys the
* information that the reply never arrived.
*
* @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
* anything
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
* was queued
*/
public boolean reply(String session, String content) {
if (content == null || content.isBlank()) {
throw new IllegalArgumentException("content is required");
}
if (rendezvous.resolve(session, content)) {
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
@@ -427,38 +427,10 @@ public final class FleetApp {
if (!allow(ctx, routeAction("GET /profiles"), null)) {
return;
}
Map<String, Object> body = new LinkedHashMap<>();
body.put("profiles", workers.profiles());
body.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
Map<String, Object> quarantined = new LinkedHashMap<>();
Map<String, Object> coolingOff = new LinkedHashMap<>();
for (String profile : workers.profiles()) {
String credentialId = quarantine.credentialIdFor().apply(profile);
if (credentialId != null) {
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
Map<String, Object> row = new LinkedHashMap<>();
row.put("credentialId", credentialId);
row.put("quarantinedForSeconds", remaining);
quarantined.put(profile, row);
});
}
String outageCredentialId = outage.credentialIdFor().apply(profile);
if (outageCredentialId != null) {
outage.outagePolicy().remainingCoolOffSeconds(outageCredentialId).ifPresent(remaining -> {
Map<String, Object> row = new LinkedHashMap<>();
row.put("credentialId", outageCredentialId);
row.put("coolingOffForSeconds", remaining);
coolingOff.put(profile, row);
});
}
}
if (!quarantined.isEmpty()) {
body.put("quarantined", quarantined);
}
if (!coolingOff.isEmpty()) {
body.put("coolingOff", coolingOff);
}
ctx.status(200).json(body);
// fleetd #297: ONE body builder, shared with fleet_profiles. Handing both doors the same
// QuarantineSource/OutageSource instances stops them reading different facts; rendering
// through the same method stops them reporting those facts differently. Both are needed.
ctx.status(200).json(FleetMcp.profilesView(workers, quarantine, outage));
}
/**
@@ -712,7 +684,19 @@ public final class FleetApp {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
messages.reply(id, content);
// fleetd #302: content is required. `.path("content").asText("")` above turns a missing key
// into "" rather than throwing, so without this check an empty/blank reply used to reach
// messages.reply(...) and silently resolve the lead's waiter — the same class of bug as the
// sibling "content is required" guards on sendMessage/askMessage below, except this one wrote
// 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
// {error, detail} envelope this file uses everywhere else.
try {
messages.reply(id, content);
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", e.getMessage()));
return;
}
ctx.status(200).json(Map.of("sessionId", id, "delivered", true));
}
@@ -864,10 +864,20 @@ public final class SessionManager implements TurnListener {
}
/**
* Gracefully drain all registered sessions on daemon shutdown. For each session that is
* {@code BUSY}, poll up to {@code timeoutNanos} for it to leave {@code BUSY}, then release it
* regardless. Non-busy sessions are released immediately. A failure releasing one session is
* logged and does not abort the rest.
* Gracefully drain all registered sessions on daemon shutdown. Non-busy sessions are released
* immediately; a {@code BUSY} one is polled until it leaves {@code BUSY}, then released
* regardless. A failure releasing one session is logged and does not abort the rest.
*
* <p>{@code timeoutNanos} is a budget for the WHOLE drain, not a grace period per session: the
* deadline is taken once, before the loop. So the first BUSY session can spend all of it, and a
* later BUSY one is then released with no wait at all. That is deliberate. This drain is only
* one phase of shutdown — {@code Fleetd} closes the message service, the push loop, the
* heartbeat, MCP and the router after it — and the whole sequence has to finish inside
* launchd's exit window. A per-session grace would let N busy members drain for N * the
* timeout, overrun that window, and get the daemon SIGKILLed part-way through; the members not
* yet reached would then get no clean release, no preserved-worktree log, and no snapshot.
* Cutting one turn short is the cheaper failure, and it is not silent: an abandoned BUSY
* session is preserved, snapshotted, and logged at WARN by {@code logPreservedForShutdown}.
*
* <p>CB-544: this is a {@link ReleaseCause#SHUTDOWN} release — the worker's pane is stopped
* (the process must end) but its worktree is preserved and its path logged. Shutdown is never
@@ -40,6 +40,7 @@ public final class FakeHerdr implements HerdrClient {
private final Map<String, List<String>> extraTabs = new LinkedHashMap<>();
private int agentNameTakenFor = 0;
private int agentPaneBusyFor = 0;
private String agentStartErrorCode = null;
private int workerTabPaneCount = 1;
private String paneCloseErrorCode = null;
private final Map<String, String> paneCloseErrorCodeFor = new ConcurrentHashMap<>();
@@ -84,6 +85,12 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/** Make every {@code agent.start} call fail with this herdr error code. */
public FakeHerdr agentStartFailsWith(String code) {
this.agentStartErrorCode = code;
return this;
}
/** Make the worker tab (w9:t2) report this many panes in {@code tab.list} (default 1). */
public FakeHerdr withWorkerTabPaneCount(int n) {
this.workerTabPaneCount = n;
@@ -311,6 +318,10 @@ public final class FakeHerdr implements HerdrClient {
+ required + "`", "invalid_request", null);
}
}
if (agentStartErrorCode != null) {
throw new HerdrException("herdr error [" + agentStartErrorCode + "]: agent.start failed",
agentStartErrorCode, null);
}
long starts = calls.stream().filter(c -> c.method().equals("agent.start")).count();
if (starts <= agentPaneBusyFor) {
throw new HerdrException(
@@ -1165,7 +1165,8 @@ class ClaudeCodeLauncherTest {
@Test
void spawnLetsAnUnrelatedHerdrErrorPropagateUnchanged() {
// Fix 1 must only special-case a "*_not_found" answer. Any other herdr failure keeps
// propagating as-is — this gate does not know how to recover from it.
// propagating as-is — this gate does not know how to recover from it. The pane still needs
// closing because spawn throws before it can return the pane id to a caller that could stop it.
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("unknown");
herdr.agentGetFailsWithAfter(0, "internal_error");
@@ -1182,8 +1183,27 @@ class ClaudeCodeLauncherTest {
() -> svc.spawn(new SpawnRequest(null, null, null)));
assertEquals("internal_error", ex.code());
assertEquals(0, paneCloseCount(herdr, "w9:pRoot_1"),
"an error this gate does not recognize is not this gate's teardown to run");
assertEquals(1, paneCloseCount(herdr, "w9:pRoot_1"),
"the unchanged error leaves spawn without a pane id, so this gate closes its orphaned pane");
}
@Test
void panePlacementClosesTheSplitPaneWhenAgentStartFails() {
FakeHerdr herdr = new FakeHerdr().agentStartFailsWith("internal_error");
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "pane", "fleetd-workers", "w #{n}", null, null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
dev.ltms.fleet.herdr.HerdrException ex = assertThrows(
dev.ltms.fleet.herdr.HerdrException.class,
() -> svc.spawn(new SpawnRequest(null, null, null)));
assertEquals("internal_error", ex.code(), "agent.start failure propagates unchanged");
assertEquals(1, paneCloseCount(herdr, "w1:pSplit"),
"the pane split for a peer that never starts is closed instead of left orphaned");
}
// --- fleetd #176 fix 2: corroborated UNKNOWN refinement --------------------------------------
@@ -154,7 +154,7 @@ class AmqpReplyInboxContractTest {
}
@Test
void releaseCancelsConsumerAndClearsHeld() throws Exception {
void releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery() throws Exception {
String target = "worker-release-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
inbox.own(target);
@@ -164,6 +164,22 @@ class AmqpReplyInboxContractTest {
inbox.release(target);
assertTrue(inbox.peek(target).isEmpty(),
"release clears the local held snapshot");
// fleetd #298: release() must not just drop the local record — the broker delivery was
// never acked, so cancelling the consumer alone leaves it unacked-but-orphaned on the
// still-open channel unless release() nacks it back with requeue=true. Prove the message
// is genuinely recoverable, not merely absent from peek: re-own the same target and
// confirm the broker redelivers it to the fresh consumer.
inbox.own(target);
List<ReplyInbox.InboxMessage> recovered = awaitPeek(inbox, target);
assertEquals(1, recovered.size(),
"a reply held (but undrained) at release() time must still be recoverable — "
+ "release() must requeue it, not silently drop it while the broker still "
+ "considers it outstanding");
assertEquals("m1", recovered.getFirst().msgId());
assertEquals("release me", recovered.getFirst().content());
inbox.ack(target, "m1");
}
}
@@ -532,6 +532,53 @@ class FleetAppTest {
assertEquals("orphan", body.get("replies").get(0).get("content").asText());
}
@Test
void replyWithMissingContentIsRejectedAndDoesNotResolveTheWaiter() throws Exception {
// fleetd #302: `.path("content").asText("")` used to turn a missing "content" key into an
// empty string that reached rendezvous.resolve, silently completing the lead's blocking wait
// with nothing. Prove the fix two ways: the bad call is rejected with 400, AND the real send
// it would have wrongly resolved is still open afterwards — a real reply completes it.
FakeHerdr herdr = new FakeHerdr().agentStatus("idle");
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
var send = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
try { return postMessage(port, "{\"content\":\"review this\",\"timeoutMs\":4000}"); }
catch (Exception e) { throw new RuntimeException(e); }
});
Thread.sleep(200); // let the background send open its rendezvous waiter
HttpResponse<String> badReply = postJson(port, "/sessions/term_a/reply", "{}");
assertEquals(400, badReply.statusCode());
JsonNode err = mapper.readTree(badReply.body());
assertEquals("bad_request", err.get("error").asText());
assertTrue(err.has("detail"));
// The waiter must still be open — a real reply now completes the ORIGINAL send.
HttpResponse<String> goodReply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}");
assertEquals(200, goodReply.statusCode());
HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS);
assertEquals(200, res.statusCode());
assertEquals("LGTM ship it", mapper.readTree(res.body()).get("reply").asText());
}
@Test
void replyWithEmptyOrWhitespaceContentIsRejectedSameAsMissing() throws Exception {
// fleetd #302 sibling case: present-but-blank content is treated the same as a missing key —
// FleetMcp's own required-content guard (fleet_reply's "content is required") makes no
// distinction between the two either, so diverging here would be a new asymmetry.
int port = startHealthy();
HttpResponse<String> empty = postJson(port, "/sessions/term_a/reply", "{\"content\":\"\"}");
assertEquals(400, empty.statusCode());
assertEquals("bad_request", mapper.readTree(empty.body()).get("error").asText());
HttpResponse<String> whitespace = postJson(port, "/sessions/term_a/reply", "{\"content\":\" \"}");
assertEquals(400, whitespace.statusCode());
assertEquals("bad_request", mapper.readTree(whitespace.body()).get("error").asText());
}
@Test
void drainRepliesReturnsEmptyForNoReplies() throws Exception {
int port = startHealthy();