Compare commits

...

14 Commits

Author SHA1 Message Date
Dai Ha a49671ceb9 #310: prevent idle reap from stopping delivered workers
CI / contract (pull_request) Successful in 1m19s
CI / build (pull_request) Successful in 2m1s
2026-09-04 13:17:51 +07:00
Dai Ha 76672ff016 #306: the post-turn phase gets the same unknown-stall escape as a turn
CI / build (push) Successful in 1m43s
CI / contract (push) Successful in 1m50s
Four latches gate delivery in Injector, and only awaitingCompletion had a way
out of a sustained unknown streak. CB-109 added that escape because a worker
stuck in a state herdr cannot classify never produces a working->idle
boundary. The same is true during post-turn housekeeping, but the escape was
never extended there.

awaitingPostTurnPickup and postTurnObserved are both released only on an
injectable sample, so a worker that goes unknown and stays there wedges: the
target is polled forever, every later message to it is blocked by the delivery
gate, and no onTurnFailed fires, so the session sits at DONE and looks healthy.
The counter did not even increment, since ++unknownSinceTurn sits inside the
awaitingCompletion short-circuit.

postTurnPending needs no escape; it is cleared on the line after the listener
call that sets it.

The escape does not set turnFailed. The delegated turn already completed and
its waiter already resolved — what is outstanding is the /clear. Failing the
turn would drive SessionManager.onFailed on a session that genuinely finished.

Not reachable in the live configuration: the path needs lifecycle.clearAfterTurn,
which fleetd.yaml does not set. It becomes reachable as soon as anyone turns
that supported knob on.

Fixes #306
2026-09-04 13:06:29 +07:00
Dai Ha 9379f92c23 #305: one definition of loopback, so a worker cannot become the primary
CI / contract (push) Successful in 52s
CI / build (push) Successful in 1m40s
ConnectionIdentity and CallerResolver each kept their own isLoopback. They
drifted: the identity resolver accepted only 127.0.0.1, the authorization
check accepted all of 127.0.0.0/8.

A caller from 127.0.0.2 therefore had its identity resolution skipped, so it
carried no terminal, and CallerResolver reads a missing terminal as "not a
worker" — which under loopback-trust, the default mode, is the primary. A
worker got spawn, stop, send and drain. The skip also happens before the PID
ancestry walk, so that defence is bypassed too.

Being strict in ConnectionIdentity was not the safe direction. That predicate
decides whether identity is resolved at all, and resolution is what demotes a
worker, so every address it excluded was one where a worker became the lead.

Measured, not assumed: on Linux the whole 127.0.0.0/8 is bound to lo, and
binding a source of 127.0.0.2 on the fleet host succeeds (curl rc=7, the
connect refused rather than the bind). On macOS the source bind fails (rc=45),
so this workstation was never exposed.

The shared predicate also accepts the IPv4-mapped IPv6 form, which neither
copy handled. That one failed in the safe direction: a primary on
::ffff:127.0.0.1 was refused as anonymous.

No transport-level test binds a real 127.0.0.2 source — it cannot run on
macOS. The reasoning is recorded on the issue.

Fixes #305
2026-09-04 12:58:21 +07:00
Dai Ha 21ff63b11d #304: the member routes report a herdr failure instead of a bare 500
CI / build (push) Successful in 1m34s
CI / contract (push) Successful in 1m53s
POST /members and DELETE /members/{paneId} were the two routes in FleetApp
with no catch (HerdrException). FleetApp has no Javalin exception mapper, so
the exception escaped as the default 500 with the body "Server Error" — no
herdr code, no herdr message. fleet_spawn and fleet_stop catch the same
exception and report a named error, so this was the same one-door-guarded
shape as #297.

Both now go through the existing herdrError mapper: 404 when herdr says the
target is gone, 502 otherwise. That is an answer the caller can act on.

stopMember matters more than spawnMember. SessionManager.release deregisters
the session, notifies the release listener and preserves a dirty worktree
before it calls launcher.stop, so a throw from that stop arrives after the
teardown the caller asked for has already happened. A bare 500 told the caller
to retry and carried nothing to explain what went wrong.

The two existing tests that asserted 500 now assert 502 and check the error
body. Neither was about the status code: one guards that a failed teardown is
not reported as a successful 204, the other that a failed spawn still closes
its tab. Both properties are unchanged.

Fixes #304
2026-09-04 12:48:21 +07:00
Dai Ha 21844b54d7 #302: fleet_reply refuses blank content too, matching fleet_send
CI / contract (push) Successful in 1m15s
CI / build (push) Successful in 3m34s
MessageService.reply now throws on blank content. FleetMcp.reply guarded only
against null, and its handler is a bare BiFunction with no try/catch, so a
whitespace-only fleet_reply left the handler as an uncaught
IllegalArgumentException instead of the clean tool error null already got.
fleet_send has always used isBlank here; reply now matches it.

The worker found that null/isBlank difference and reported it as a correction
to my ticket, which had quoted the guard wrongly. It was right: I grepped the
error string and assumed the condition matched its sibling.
2026-09-04 12:36:23 +07:00
Dai Ha 4769481515 Merge #302: a reply with no content is refused, not silently delivered
REST read content with .path("content").asText(""), so a body missing the key
became an empty string that resolved the lead's waiter. The turn completed and
the lead saw a member that finished and reported nothing, indistinguishable
from one that genuinely said nothing. The guard went into MessageService.reply,
which both doors call, rather than being written a second time in FleetApp.
2026-09-04 12:33:47 +07:00
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 85c90d440a #297: map HerdrException on GET /agents and /members; GET /profiles reports quarantine + cool-off
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m33s
Two REST-only visibility gaps, both against the same shared instances FleetMcp
reads (BackendQuarantine/BackendOutagePolicy), never recomputed:

- GET /agents and GET /members let a HerdrException escape uncaught, outside
  the {error, detail} envelope every other failure path in FleetApp uses.
  Both now route through the existing herdrError() helper, matching healthz/
  sessionStatus. GET /members is the endpoint's own comment names as the
  out-of-band path a lead falls back to when its MCP mount drops.
- GET /profiles omitted the two outage states fleet_profiles already reports:
  quarantined (CB-578 stage B) and coolingOff (fleetd #201 Unit 5). FleetApp
  now takes the SAME FleetMcp.QuarantineSource/OutageSource instances Fleetd
  wires into FleetMcp (extracted to local vars in Fleetd.java so both doors
  share one object, not two independently-built copies of the same rule).

FleetMcp itself is unchanged. Item 3 of the ticket (a capacity block on
GET /members) is explicitly out of scope and was not added.
2026-09-04 12:12:35 +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
16 changed files with 659 additions and 76 deletions
@@ -625,6 +625,19 @@ public final class Fleetd {
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
}
// fleetd #297: named once and reused verbatim below for FleetApp's GET /profiles, rather than
// built a second time — two independently-constructed sources reading the SAME BackendQuarantine
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
// closure in only one of the two places), exactly the shape #284 was.
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine);
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy);
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, metrics, new FleetMcp.CapacitySource(profile -> liveCountRef.get().apply(profile),
profile -> {
@@ -636,15 +649,9 @@ public final class Fleetd {
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
health != null && health.notifications() != null && health.notifications().configured());
}),
new FleetMcp.QuarantineSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine),
quarantineSource,
leadMailbox,
new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy),
outageSource,
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
@@ -715,9 +722,12 @@ public final class Fleetd {
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
// fleetd #111: live (re-read-per-request) memberCredentials view for GET /member-credentials —
// same hot-reload shape as the memberCredentials supplier passed to ClaudeCodeLauncher above.
// fleetd #297: quarantineSource/outageSource are the SAME instances passed to FleetMcp above —
// GET /profiles must report the identical quarantine/cool-off facts as fleet_profiles.
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics, deliverable,
() -> MemberCredentialPolicyView.of(config.get().memberCredentials())).build();
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
quarantineSource, outageSource).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("fleetd listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
@@ -262,11 +262,15 @@ public final class CallerResolver {
return token.isEmpty() ? null : token;
}
/**
* fleetd #305: delegates to {@link ConnectionIdentity#isLoopback}. This used to be a second,
* independent copy of the same rule, and the two drifted: this one accepted all of
* {@code 127.0.0.0/8}, {@code ConnectionIdentity}'s accepted only {@code 127.0.0.1}. A caller
* from {@code 127.0.0.2} therefore had its identity skipped (so it had no terminal) and was
* then read as loopback here — which under loopback-trust is the primary. Sharing the inputs
* would not have prevented that; only sharing the computation does.
*/
private static boolean isLoopback(String remoteAddr) {
if (remoteAddr == null) {
return false;
}
return remoteAddr.equals("127.0.0.1") || remoteAddr.equals("::1")
|| remoteAddr.equals("0:0:0:0:0:0:0:1") || remoteAddr.startsWith("127.");
return ConnectionIdentity.isLoopback(remoteAddr);
}
}
@@ -157,6 +157,7 @@ public final class Injector {
boolean awaitingCompletion; // a delivered message's turn is not yet known-complete
boolean turnObserved; // saw a real `working` sample since that delivery (turn ran)
int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109)
int unknownSincePostTurn; // the same, for the post-turn housekeeping phase (fleetd #306)
int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114)
boolean postTurnPending; // completion observed; adapter housekeeping has not started yet
boolean awaitingPostTurnPickup;
@@ -217,10 +218,12 @@ public final class Injector {
t.awaitingPickup = false;
t.injectableSincePickup = 0;
t.unknownSinceTurn = 0;
t.unknownSincePostTurn = 0;
t.notReadySincePoll = 0;
if (t.awaitingCompletion) t.turnObserved = true;
} else if (status.injectable()) { // IDLE or BLOCKED
t.unknownSinceTurn = 0;
t.unknownSincePostTurn = 0;
if (t.awaitingPostTurnPickup) {
if (++t.injectableSincePostTurnPickup >= PICKUP_GRACE_POLLS) {
t.awaitingPostTurnPickup = false;
@@ -315,6 +318,24 @@ public final class Injector {
t.unknownSinceTurn = 0;
turnFailed = true;
}
// fleetd #306: the same escape for the post-turn housekeeping phase. Four latches
// gate delivery (awaitingCompletion, postTurnPending, awaitingPostTurnPickup,
// postTurnObserved) and only the first had a way out of a sustained unknown streak —
// a gate that closed one direction only. The other two below are released here as
// well; postTurnPending needs no escape because it is cleared unconditionally on the
// line after the listener call that sets it.
//
// This does NOT set turnFailed. The delegated turn already completed and its waiter
// already resolved — what is outstanding is adapter housekeeping (the `/clear`).
// Reporting a turn failure here would drive SessionManager.onFailed on a session
// that genuinely finished its work, which is a worse lie than the wedge.
if ((t.awaitingPostTurnPickup || t.postTurnObserved)
&& ++t.unknownSincePostTurn >= TURN_STALL_GRACE_POLLS) {
t.awaitingPostTurnPickup = false;
t.postTurnObserved = false;
t.injectableSincePostTurnPickup = 0;
t.unknownSincePostTurn = 0;
}
}
// Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or
@@ -60,7 +60,30 @@ public final class ConnectionIdentity {
return pid > 0 ? cwds.cwdForPid(pid) : null;
}
private static boolean isLoopback(String addr) {
return "127.0.0.1".equals(addr) || "::1".equals(addr) || "0:0:0:0:0:0:0:1".equals(addr);
/**
* Whether {@code addr} is a same-host address, and therefore one whose peer PID is worth
* looking up. <strong>This is the one definition of loopback in the daemon</strong> —
* {@code CallerResolver} calls it rather than keeping its own, because the two used to differ
* and that difference was a privilege escalation (fleetd #305).
*
* <p>The whole of {@code 127.0.0.0/8} counts, not just {@code 127.0.0.1}. On Linux every
* address in that range is bound to {@code lo} by default, so a process can connect to
* {@code 127.0.0.1:8765} with a source address of {@code 127.0.0.2} — measured on the Linux
* fleet host, where binding that source succeeds.
*
* <p><strong>Being strict here does not make the daemon safer; it makes it unsafe.</strong>
* That reads backwards, so it is worth stating plainly. This predicate does not decide whether
* a caller is trusted — it decides whether the caller's identity is <em>resolved at all</em>.
* Returning false means {@link #resolve} answers "no terminal", and downstream a caller with no
* terminal is treated as the primary under loopback-trust. So every address excluded here is an
* address on which a worker silently becomes the lead. Widening a check normally weakens it;
* widening this one is what closes the hole.
*/
public static boolean isLoopback(String addr) {
if (addr == null) {
return false;
}
String a = addr.startsWith("::ffff:") ? addr.substring(7) : addr; // IPv4-mapped IPv6
return a.startsWith("127.") || "::1".equals(a) || "0:0:0:0:0:0:0:1".equals(a);
}
}
@@ -813,7 +813,12 @@ public final class FleetMcp {
return error("fleet_reply is for workers only — could not identify the calling worker "
+ "from the connection");
}
if (content == null) {
// fleetd #302: isBlank, not == null, to match fleet_send's own guard above. MessageService
// .reply now REJECTS blank content, and this handler is a bare BiFunction with no try/catch
// around it — so a whitespace-only fleet_reply would leave here as an uncaught
// IllegalArgumentException instead of this clean tool error. Null and whitespace are the
// same mistake by the caller and must get the same answer.
if (isBlank(content)) {
return error("content is required");
}
messages.reply(callerTerminal, content);
@@ -1015,6 +1020,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 +1070,7 @@ public final class FleetMcp {
if (!coolingOff.isEmpty()) {
result.put("coolingOff", coolingOff);
}
return text(json(result));
return result;
}
/**
@@ -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
@@ -9,6 +9,7 @@ import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
@@ -86,6 +87,12 @@ public final class FleetApp {
// absent() (the honest "no policy configured" view) for every constructor that does not wire
// a real one, so existing legacy call sites keep building without knowing this field exists.
private final Supplier<MemberCredentialPolicyView> memberCredentials;
// fleetd #297: the SAME shared sources FleetMcp.profiles/fleet_profiles reads (BackendQuarantine
// and BackendOutagePolicy are each one instance for the whole daemon — see Fleetd wiring) so
// GET /profiles cannot drift from fleet_profiles about which profile is quarantined/cooling off.
// .none() (the honest "feature not wired" view) for every constructor that does not pass one.
private final FleetMcp.QuarantineSource quarantine;
private final FleetMcp.OutageSource outage;
private final ObjectMapper mapper = new ObjectMapper();
/**
@@ -146,6 +153,23 @@ public final class FleetApp {
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials) {
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics,
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
}
/**
* @param quarantine the SAME {@link FleetMcp.QuarantineSource} instance passed to {@code
* FleetMcp} (fleetd #297), so {@code GET /profiles} reports the identical
* exhaustion-quarantine facts as {@code fleet_profiles} rather than a second,
* independently-computed copy
* @param outage the SAME {@link FleetMcp.OutageSource} instance passed to {@code FleetMcp} —
* see {@code quarantine}; a SEPARATE check from it, never merged in
*/
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
this.herdr = herdr;
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
this.workers = workers;
@@ -156,6 +180,8 @@ public final class FleetApp {
this.auth = auth;
this.metrics = metrics;
this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent;
this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none();
this.outage = outage != null ? outage : FleetMcp.OutageSource.none();
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
@@ -338,8 +364,15 @@ public final class FleetApp {
if (!allow(ctx, routeAction("GET /agents"), null)) {
return;
}
ctx.status(200).json(Map.of("agents",
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
try {
ctx.status(200).json(Map.of("agents",
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
} catch (HerdrException e) {
// fleetd #297: workers.list() reaches herdr — a transport failure must land in the same
// {error, detail} envelope every other failure path here uses, not escape as a bare
// exception and leave Javalin's default handling to respond outside the JSON contract.
herdrError(ctx, e);
}
}
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
@@ -347,40 +380,57 @@ public final class FleetApp {
if (!allow(ctx, routeAction("GET /members"), null)) {
return;
}
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
List<Map<String, Object>> out = sessions.rosterResolved().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
.toList();
Map<String, Object> body = new LinkedHashMap<>();
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
// drops reads this endpoint. Drop the alias once nothing reads it.
body.put("members", out);
body.put("workers", out);
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
Map.of("count", st.count(), "costBytes", st.costBytes())));
ctx.status(200).json(body);
try {
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
List<Map<String, Object>> out = sessions.rosterResolved().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
.toList();
Map<String, Object> body = new LinkedHashMap<>();
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
// drops reads this endpoint. Drop the alias once nothing reads it.
body.put("members", out);
body.put("workers", out);
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
Map.of("count", st.count(), "costBytes", st.costBytes())));
ctx.status(200).json(body);
} catch (HerdrException e) {
// fleetd #297: same reasoning as agents() above — this is the out-of-band roster a lead
// falls back to when its MCP mount drops, so it must stay inside the JSON error contract
// exactly when herdr is briefly unreachable, not escape as a bare exception.
herdrError(ctx, e);
}
}
/** The configured worker profiles and which one a no-argument spawn uses. */
/**
* The configured worker profiles, which one a no-argument spawn uses, and (fleetd #297) the two
* outage states {@code fleet_profiles} already reports: {@code quarantined} (CB-578 stage B —
* the backend reported it out of capacity) and {@code coolingOff} (fleetd #201 Unit 5 — the
* credential threw repeated non-exhaustion backend errors). Both are read from the SAME shared
* {@link FleetMcp.QuarantineSource}/{@link FleetMcp.OutageSource} instances {@code FleetMcp}
* reads, never recomputed, so the two doors cannot disagree about which profile is down and why.
* Independent checks, so a profile can appear in both maps at once; each map is present only
* when at least one profile is in that state.
*/
private void profiles(Context ctx) {
if (!allow(ctx, routeAction("GET /profiles"), null)) {
return;
}
ctx.status(200).json(Map.of(
"profiles", workers.profiles(),
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile()));
// 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));
}
/**
@@ -460,6 +510,12 @@ public final class FleetApp {
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
} catch (PeerUnreachableException e) {
ctx.status(502).json(Map.of("error", "spawn_timeout", "detail", e.getMessage()));
} catch (HerdrException e) {
// fleetd #304: not every herdr failure on the spawn path is a readiness timeout, so
// PeerUnreachableException above does not cover this. Without this catch the exception
// escapes to Javalin's default 500, while fleet_spawn reports the same failure as a
// clean named error (FleetMcp.spawn) — the #297 one-door-guarded shape.
herdrError(ctx, e);
}
}
@@ -480,13 +536,29 @@ public final class FleetApp {
return (s == null || s.isBlank()) ? null : s;
}
/** Tear a worker down by pane id. */
/**
* Tear a worker down by pane id.
*
* <p>fleetd #304: the {@code HerdrException} catch is not cosmetic. {@code release} deregisters
* the session, notifies the release listener and preserves a dirty worktree <em>before</em> it
* calls {@code launcher.stop}, so a throw from that stop arrives after the teardown the caller
* asked for has already happened. Letting it escape gave Javalin's default 500, which tells the
* caller to retry — and the retry finds nothing in the registry, reaches the same stop, and
* throws again, so it can never succeed. {@code herdrError} instead answers 404 ("the pane is
* gone, stop retrying") or 502 ("herdr is upstream and broken, a retry may help"), matching what
* {@code fleet_stop} reports for the same failure.
*/
private void stopMember(Context ctx) {
String paneId = ctx.pathParam("paneId");
if (!allow(ctx, routeAction("DELETE /members/{paneId}"), paneId)) {
return;
}
sessions.release(paneId);
try {
sessions.release(paneId);
} catch (HerdrException e) {
herdrError(ctx, e);
return;
}
ctx.status(204);
}
@@ -634,7 +706,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));
}
@@ -61,6 +61,8 @@ public final class SessionManager implements TurnListener {
private final LongSupplier nowNanos;
private final int contextCap;
private final boolean clearAfterTurn;
/** Null in production; test seam for the interval before an idle session's conditional release. */
private final Consumer<MemberSession> beforeIdleRelease;
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
/**
* CB-586: the repo root the fleet actually works in, remembered the first time a worktree
@@ -102,13 +104,23 @@ public final class SessionManager implements TurnListener {
}
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
int contextCap, boolean clearAfterTurn) {
int contextCap, boolean clearAfterTurn) {
this(launcher, worktrees, nowNanos, contextCap, clearAfterTurn, null);
}
/**
* Package-private constructor for a deterministic reap/delivery race test. Production callers
* use the constructor above, whose null hook adds no callback or lock to an ordinary reap.
*/
SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
int contextCap, boolean clearAfterTurn, Consumer<MemberSession> beforeIdleRelease) {
this.launcher = launcher;
this.worktrees = worktrees;
this.presence = new PresenceFleet(this);
this.nowNanos = nowNanos;
this.contextCap = contextCap;
this.clearAfterTurn = clearAfterTurn;
this.beforeIdleRelease = beforeIdleRelease;
}
/**
@@ -272,10 +284,27 @@ public final class SessionManager implements TurnListener {
*/
private void release(String paneId, ReleaseCause cause) {
MemberSession removed = registry.remove(paneId);
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
}
/**
* Tear a session down only while {@code expected} is still its registry value. A lifecycle
* transition replaces the immutable record, so this prevents a reap based on an old READY or
* DONE record from stopping a worker that delivery has made BUSY.
*/
private boolean releaseIfCurrent(MemberSession expected, ReleaseCause cause) {
if (!registry.remove(expected.paneId(), expected)) {
return false;
}
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
return true;
}
private void releaseRemoved(String paneId, MemberSession removed, PeerHandle removedHandle,
ReleaseCause cause) {
// fleetd #209: remove right alongside the registry entry so a released session's handle is
// never leaked — but keep the local reference below, so the id can still be resolved for
// the ReleaseDetail this teardown notifies with.
PeerHandle removedHandle = handles.remove(paneId);
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
String snapshotRef = null;
if (removed != null) {
@@ -846,14 +875,18 @@ public final class SessionManager implements TurnListener {
}
long idleNanos = now - s.lastActivityAtNanos();
if (idleNanos > idleTtlNanos) {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
// CB-581: one session that fails to release must not abort the whole reaping pass —
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
try {
release(s.paneId());
reaped++;
if (beforeIdleRelease != null) {
beforeIdleRelease.accept(s);
}
if (releaseIfCurrent(s, ReleaseCause.COMPLETED)) {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
reaped++;
}
} catch (RuntimeException e) {
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
@@ -864,10 +897,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
@@ -413,4 +413,29 @@ class CallerResolverTest {
assertThrows(IllegalArgumentException.class, () -> new CallerResolver(id, true, null));
assertThrows(IllegalArgumentException.class, () -> new CallerResolver(id, true, " "));
}
@Test
void aWorkerOnAnyLoopbackSourceAddressIsStillAWorkerNotThePrimary() {
// fleetd #305: the escalation. ConnectionIdentity used to accept only 127.0.0.1, so a
// worker connecting from 127.0.0.2 resolved to no terminal, and this resolver's own
// (wider) loopback check then made it the PRIMARY — granting spawn, stop, send and drain.
// Measured on the Linux fleet host: binding a source of 127.0.0.2 succeeds there, so the
// path is real and not theoretical.
CallerResolver r = new CallerResolver(workerIdentity(), false, null);
for (String src : new String[]{"127.0.0.1", "127.0.0.2", "127.1.2.3", "::ffff:127.0.0.2"}) {
Principal p = r.resolve(src, 55555, null);
assertEquals(Role.WORKER, p.role(), "a worker must stay a worker from source " + src);
assertEquals("term_a", p.terminal(), "worker terminal from source " + src);
}
}
@Test
void aNonWorkerOnAnyLoopbackSourceAddressIsStillThePrimary() {
// The other direction of the same fix: widening the identity check must not demote a
// legitimate same-host primary that happens to connect from another 127.* address.
CallerResolver r = new CallerResolver(nonWorkerIdentity(), false, null);
for (String src : new String[]{"127.0.0.1", "127.0.0.2", "::ffff:127.0.0.1"}) {
assertEquals(Role.PRIMARY, r.resolve(src, 55555, null).role(), "source " + src);
}
}
}
@@ -297,6 +297,68 @@ class InjectorTest {
assertTrue(inj.activeTargets().isEmpty(), "the wedged target is reclaimed, not polled forever");
}
/** A listener whose post-turn housekeeping always starts, as SessionManager's does with clearAfterTurn on. */
private static final class PostTurnListener implements TurnListener {
@Override public void onTurnComplete(String target) { }
@Override public boolean hasPostTurnAction(String target) { return true; }
@Override public boolean onTurnCompleteWithPostAction(String target) { return true; }
}
@Test
void aWorkerThatWedgesInUnknownAwaitingPostTurnPickupIsReleased() {
// fleetd #306: the post-turn phase had no way out of a sustained unknown streak, so the
// pickup latch stayed set, the target was polled forever, and every later message to it was
// blocked by the delivery gate — while the session still looked healthy.
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver
inj.onStatus(T, AgentStatus.WORKING); // turn starts
inj.onStatus(T, AgentStatus.IDLE); // turn completes; housekeeping dispatched
for (int i = 0; i < STALL_SAMPLES; i++) inj.onStatus(T, AgentStatus.UNKNOWN); // then wedges
assertTrue(inj.activeTargets().isEmpty(),
"a target wedged awaiting post-turn pickup must be reclaimed, not polled forever");
assertEquals(List.of(), cap.failed,
"the delegated turn already completed — a stuck /clear must not be reported as a failed turn");
}
@Test
void aWorkerThatWedgesInUnknownAfterPickingUpTheResetIsReleased() {
// The sibling latch. postTurnObserved is set when the reset is seen picked up (WORKING) and
// is cleared only on a later injectable sample, so a wedge right after pickup sticks too.
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
inj.onStatus(T, AgentStatus.WORKING);
inj.onStatus(T, AgentStatus.IDLE); // turn complete; reset dispatched
inj.onStatus(T, AgentStatus.WORKING); // reset picked up -> postTurnObserved
for (int i = 0; i < STALL_SAMPLES; i++) inj.onStatus(T, AgentStatus.UNKNOWN);
assertTrue(inj.activeTargets().isEmpty(), "a wedge after reset pickup must also be reclaimed");
assertEquals(List.of(), cap.failed, "still not a turn failure");
}
@Test
void aBriefUnknownDuringPostTurnHousekeepingDoesNotDropTheLatch() {
// The other direction: the escape must not fire on a glitch, or the queued next delegation
// would overtake housekeeping that is still running.
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
inj.enqueue(T, "first", TestTurnTokens.inert(T));
inj.enqueue(T, "second", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
inj.onStatus(T, AgentStatus.WORKING);
inj.onStatus(T, AgentStatus.IDLE); // first completes; reset dispatched
for (int i = 0; i < 10; i++) inj.onStatus(T, AgentStatus.UNKNOWN); // well under the grace
assertFalse(inj.activeTargets().isEmpty(), "a brief glitch must not release the post-turn latch");
assertEquals(List.of("first"), sent(), "the queued delegation must not overtake housekeeping");
}
@Test
void aTransientUnknownGlitchNeitherFailsNorBlocksCompletion() {
Captor cap = new Captor();
@@ -20,6 +20,17 @@ class ConnectionIdentityTest {
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.0.0.1", 55555));
}
@Test
void resolvesWorkerFromAnyLoopbackSourceAddressNotJust127001() {
// fleetd #305. On Linux the whole 127.0.0.0/8 is bound to lo, so a worker can connect with
// a source address of 127.0.0.2. If identity resolution skips that address the caller has
// no terminal, and a caller with no terminal is the primary under loopback-trust — so this
// must resolve the worker, not null.
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.0.0.2", 55555));
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.1.2.3", 55555));
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("::ffff:127.0.0.2", 55555));
}
@Test
void nullForOffHostCaller() {
// A non-loopback peer can't be an on-host worker → treat as primary/unknown.
@@ -326,6 +326,26 @@ class FleetMcpTest {
assertEquals("orphan", drained.getFirst().content());
}
@Test
void replyWithBlankContentIsACleanToolErrorNotAnUncaughtException() {
// fleetd #302: MessageService.reply now REJECTS blank content by throwing. fleet_reply's
// handler is a bare BiFunction with no try/catch around it, so if this guard only checked
// `== null` (as it did), a whitespace-only reply would leave the handler as an uncaught
// IllegalArgumentException instead of a tool error the caller can read. Null and whitespace
// are the same caller mistake and must get the same answer — the sibling fleet_send guard
// has always used isBlank for exactly this reason.
for (String blank : new String[] {null, "", " ", "\n\t"}) {
McpSchema.CallToolResult res = assertDoesNotThrow(
() -> FleetMcp.reply(messages, "term_a", blank),
"blank content must be refused as a tool error, never thrown out of the handler");
assertEquals(Boolean.TRUE, res.isError(), "blank content is an error result");
assertTrue(textOf(res).contains("content is required"),
"the error names the missing argument: " + textOf(res));
}
assertEquals(0, messages.drainReplies("term_a").size(),
"a refused reply must not reach the inbox");
}
@Test
void bridgePollWithTargetDrainsReplies() {
// A reply with no open send queues it in the inbox.
@@ -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");
}
}
@@ -19,6 +19,10 @@ import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.Worktrees;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.member.MemberCredentialPolicyView;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementPolicies;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
@@ -32,6 +36,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.function.Predicate;
import static org.junit.jupiter.api.Assertions.*;
@@ -70,6 +75,19 @@ class FleetAppTest {
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
Worktrees worktrees, Predicate<String> deliverable) {
return start(herdr, workerBaseUrl, allow, placement, worktrees, deliverable,
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
}
/**
* fleetd #297: same wiring as above, plus the two SAME shared sources {@code GET /profiles}
* must read — lets a test prove the quarantined/coolingOff facts it reports come from a real
* {@link dev.ltms.fleet.placement.BackendQuarantine}/{@link
* dev.ltms.fleet.placement.BackendOutagePolicy}, exactly like {@code fleet_profiles}'s own tests.
*/
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
Worktrees worktrees, Predicate<String> deliverable,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
placement, "fleet", "worker: {profile} #{n}", null, null, null);
@@ -90,8 +108,9 @@ class FleetAppTest {
// it directly so the inbox contract holds for those endpoints.
inbox.own("term_a");
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null,
null, null, id -> this.presence.isPresent(id) || deliverable.test(id))
app = new FleetApp(herdr, herdr, workers, sessions, messages, this.presence, null,
null, null, id -> this.presence.isPresent(id) || deliverable.test(id),
MemberCredentialPolicyView::absent, quarantine, outage)
.build().start("127.0.0.1", 0);
return app.port();
}
@@ -164,6 +183,22 @@ class FleetAppTest {
assertEquals("idle", agents.get(0).get("status").asText());
}
/**
* fleetd #297 gap 1: {@code workers.list()} reaches herdr, and a transport failure there must
* land in the same {@code {error, detail}} envelope every other failure path in this file uses
* (see {@code herdrError}), not escape as a bare exception outside the JSON contract.
*/
@Test
void agentsMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception {
FakeHerdr down = new FakeHerdr().healthy(false);
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "GET", "/agents");
assertEquals(502, res.statusCode(), res.body());
JsonNode body = mapper.readTree(res.body());
assertEquals("herdr_error", body.get("error").asText());
assertTrue(body.has("detail"), res.body());
}
@Test
void spawnWorkerLandsInOwnTabInWorkerSpaceAndInjectsBaseUrl() throws Exception {
FakeHerdr herdr = new FakeHerdr();
@@ -203,6 +238,42 @@ class FleetAppTest {
JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body());
assertEquals("ltms-local", body.get("default").asText());
assertEquals("ltms-local", body.get("profiles").get(0).asText());
assertFalse(body.has("quarantined"), "nothing is quarantined, so the key is omitted: " + body);
assertFalse(body.has("coolingOff"), "nothing is cooling off, so the key is omitted: " + body);
}
/**
* fleetd #297 gap 2: {@code GET /profiles} must report the same two outage states {@code
* fleet_profiles} does — CB-578 stage B exhaustion quarantine and fleetd #201 Unit 5 cool-off —
* reading the SAME shared {@link BackendQuarantine}/{@link BackendOutagePolicy} instances rather
* than recomputing them. The two checks are independent, and this profile is deliberately put in
* both states at once, matching {@code FleetMcpTest}'s own coverage of that overlap.
*/
@Test
void profilesReportsQuarantineAndCoolingOffFromTheSameSharedSources() throws Exception {
FakeHerdr herdr = new FakeHerdr();
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine("shared-openai");
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine);
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "t1", "API Error: rate limited");
outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); // 2nd distinct target starts the incident
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "tab", new GitWorktrees(),
ignored -> false, quarantineSource, outageSource);
JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body());
assertTrue(body.has("quarantined"), body.toString());
assertEquals("shared-openai",
body.get("quarantined").get("ltms-local").get("credentialId").asText());
assertEquals(1800,
body.get("quarantined").get("ltms-local").get("quarantinedForSeconds").asLong());
assertTrue(body.has("coolingOff"), body.toString());
assertEquals("shared-openai",
body.get("coolingOff").get("ltms-local").get("credentialId").asText());
assertEquals(60, body.get("coolingOff").get("ltms-local").get("coolingOffForSeconds").asLong());
}
@Test
@@ -239,6 +310,23 @@ class FleetAppTest {
"liveStatus is unknown when herdr has no matching pane");
}
/**
* fleetd #297 gap 1: same reasoning as {@code agentsMapsAHerdrFailureToTheJsonErrorEnvelope} —
* {@code GET /members} is the endpoint's own comment names as "the out-of-band path a lead falls
* back to when its MCP mount drops", so it must stay inside the {@code {error, detail}} envelope
* exactly when herdr is briefly unreachable.
*/
@Test
void membersMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception {
FakeHerdr down = new FakeHerdr().healthy(false);
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "GET", "/members");
assertEquals(502, res.statusCode(), res.body());
JsonNode body = mapper.readTree(res.body());
assertEquals("herdr_error", body.get("error").asText());
assertTrue(body.has("detail"), res.body());
}
@Test
void spawnWithACwdParamRootsTheWorkerThere() throws Exception {
FakeHerdr herdr = new FakeHerdr();
@@ -327,7 +415,11 @@ class FleetAppTest {
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(99);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
assertEquals(500, req(port, "POST", "/members").statusCode());
// fleetd #304: 502, not Javalin's default 500 — the herdr failure is named, and the body
// carries herdr's own message, matching what fleet_spawn reports for the same failure.
HttpResponse<String> res = req(port, "POST", "/members");
assertEquals(502, res.statusCode());
assertEquals("herdr_error", mapper.readTree(res.body()).get("error").asText());
assertTrue(herdr.called("tab.create"), "a tab was created before the failed start");
assertEquals("w9:t2", params(herdr, "tab.close").get("tab_id"), "orphaned tab must be closed");
}
@@ -444,6 +536,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();
@@ -598,7 +737,12 @@ class FleetAppTest {
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
// A genuine teardown failure must surface, not be reported as a successful 204.
assertEquals(500, req(port, "DELETE", "/members/w9:pW").statusCode());
// fleetd #304: it surfaces as a named 502 rather than Javalin's default 500. The property
// this test guards is "not 204" and the herdr detail reaching the caller — a bare 500 gave
// the body "Server Error" and said nothing about herdr.
HttpResponse<String> res = req(port, "DELETE", "/members/w9:pW");
assertEquals(502, res.statusCode());
assertEquals("herdr_error", mapper.readTree(res.body()).get("error").asText());
assertFalse(herdr.called("tab.close"), "tab is not removed when the pane close failed");
}
@@ -611,4 +755,5 @@ class FleetAppTest {
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
assertTrue(herdr.called("tab.close"));
}
}
@@ -178,14 +178,19 @@ class SessionManagerTest {
}
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
boolean clearAfterTurn) {
boolean clearAfterTurn) {
return sessionManager(herdr, clock, contextCap, clearAfterTurn, null);
}
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
boolean clearAfterTurn, java.util.function.Consumer<MemberSession> hook) {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn);
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn, hook);
}
@Test
@@ -670,6 +675,24 @@ class SessionManagerTest {
"BUSY session remains");
}
@Test
void reapIdleDoesNotReleaseSessionDeliveredAfterItsEligibilityCheck() {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager[] manager = new SessionManager[1];
SessionManager sessions = sessionManager(herdr, () -> clock[0], 0, false,
session -> manager[0].onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId())));
manager[0] = sessions;
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
clock[0] = 11;
assertEquals(0, sessions.reapIdle(10), "delivery replaces the idle snapshot before release");
assertEquals(MemberSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(),
"a just-delivered session stays registered and busy");
assertFalse(herdr.called("pane.close"), "the busy session pane is not stopped");
}
@Test
void doneSessionPastIdleTtlIsReaped() {
long[] clock = {0};