Compare commits

...

8 Commits

Author SHA1 Message Date
Dai Ha 8d2763e354 Add authorization entry-point audit 2026-09-04 10:32:23 +07:00
Dai Ha 66e5247b6d Merge #275 (PR #279): sweep an ASKING ticket on a definite teardown
CI / contract (push) Successful in 49s
CI / build (push) Successful in 1m39s
A member torn down while parked in fleet_ask left its async ticket pending
for good. resolveQuestion had already closed the forward waiter, so
abandon()'s waiter branch found nothing; the 'question == null' guard then
excluded the task from the matching loop. By the time the worker's own ask
lapsed (~55-115s), the released session was gone from the roster, so
nothing was left to call abandon() on that target again. fleet_poll{ticket}
reported PENDING forever.

The ticket told the worker to drop the 'question == null' guard. That was
wrong, and the worker said so with evidence: an existing test
(abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer) deliberately pins that
an ASKING ticket must SURVIVE abandon(), because the primary may be mid
answer() for that same turn. Widening the shared method would have traded
this bug for a worse one — a health guess killing a live conversation.

So the fix splits the two callers by what they actually know:

  - sessions.onRelease (fleet_stop / idle reaper) knows the pane is being
    stopped right now, so it sweeps: sweepAsking=true.
  - FleetHealthMonitor keeps sweepAsking=false. GONE/NEVER_READY is a
    classification from the live agent list, not a teardown it performed.

I verified the reachability chain myself rather than taking it on trust.
FleetHealth.decide returns GONE before it can ever return
DELEGATION_ORPHANED; FleetHealthMonitor.reportTransition returns early when
previous == next; and terminal() is GONE/NEVER_READY only. So after the one
GONE transition fires and no-ops, nothing re-fires. Every link holds.

Verified: the real merge into current main builds green (1283 tests), the
protective test still passes untouched, and my own mutation — reverting the
sweepAsking widening — fails the new test with 'a released target's open ask
can never resume, so it must fail right here'.
2026-09-04 10:18:59 +07:00
Dai Ha 1e41bd63b4 Merge #274 (PR #277): clean up the worktree when add() fails after creating it
CI / contract (push) Successful in 59s
CI / build (push) Successful in 1m34s
GitWorktrees.add() created the worktree and branch, then ran more steps that
can throw — requireCredentialFreeHttpsOrigin among them, which is an
intended security refusal, not an IO accident. Any throw meant add() never
returned, so SessionManager.acquireWithWorktree never learned the path, its
'if (path != null)' cleanup could not fire, and the worktree and branch
leaked with nothing tracking them. Every OTHER exit from that method was
cleaned up correctly; only the exits inside add() were uncounted.

add() now cleans up what it created before rethrowing, reusing remove() and
additionally deleting the branch — a branch that never finished provisioning
has no session and no PR behind it. Worktree first, since a checked-out
branch cannot be deleted. Cleanup failure is logged and never masks the
original exception.

Verified by me: the real merge into current main builds green (1281 tests),
and I reran the mutation myself without git stash — dropping the cleanup
call fails the new test with 'the worktree directory leaked after a
post-creation step threw'.

The test drives add() itself through the existing afterWorktreeAdded seam,
so the failure happens after the worktree exists rather than downstream in
another caller.
2026-09-04 10:14:05 +07:00
Dai Ha 4887d03d88 #275: abandon() sweeps an ASKING ticket only on a definite teardown
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 2m5s
Confirmed reachable: a target torn down for good (fleet_stop / the idle
reaper) while its async ticket sits in fleet_ask (Phase.ASKING) got
permanently stuck. resolveQuestion already closes the forward waiter, the
question == null guard excluded the task from abandon()'s sweep, and by the
time the worker's own fleet_ask lapses (~55-115s) the released session no
longer appears in FleetHealthMonitor's roster, so nothing ever calls
abandon() again. fleet_poll{ticket} then reports PENDING forever.

Add abandon(target, reason, sweepAsking) — sessions.onRelease (a definite
teardown: the pane is being stopped right now) passes true and now fails the
ASKING ticket and closes its reverse-rendezvous ask. FleetHealthMonitor's
health-classification call keeps the 2-arg overload (sweepAsking=false):
a GONE/NEVER_READY reading is a guess from the live agent list, not a
teardown it performed, and abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer
already covers why an active ask must survive that guess (the primary may
be mid-answer for the same turn). hasOrphanedDelegation is left unchanged
for the same reason — it must not flag a live, active ask as orphaned.

Proven with a test driving the real public sequence (sendAsync -> ask ->
abandon(..., true)), not a hand-built task map; reverted the widening to
confirm it goes red, then restored it.
2026-09-04 10:13:56 +07:00
Dai Ha 7d5434455d implementer: never git stash — the stash stack is shared across worktrees
CI / build (push) Successful in 1m41s
CI / contract (push) Successful in 1m48s
A worker's worktree is isolated; refs/stash is not. It is one stack shared
by the primary's checkout and every worker worktree of this repo.

This bit a real worker today. Two ran in parallel; one called git stash
while the other was mid-edit, and the second worker's in-progress change
was silently overwritten by the first's stashed content. It recovered by
retyping the edit and diffing to confirm, and pushed the other worker's
change back onto the stack untouched — but nothing warned either of them,
and nothing would have.

Measured before writing this: 'git stash list' from a worker worktree and
from the primary's checkout return byte-identical output, and refs/stash
is a single common ref, not a per-worktree one.

The branch already IS the isolation, so the skill now points at committing
a wip commit or writing a patch file instead.
2026-09-04 10:10:08 +07:00
Dai Ha f71ee4926e Merge #273 (PR #278): validate exhaustedPattern at config load
CI / contract (push) Successful in 44s
CI / build (push) Successful in 1m44s
A malformed exhaustedPattern passed FleetConfig.load and then crashed the
daemon at startup, in Fleetd.main's unguarded Pattern.compile, with a
message naming neither the profile nor the key. Its sibling errorPattern
had a load-time validator whose own javadoc explains exactly why that is
bad. The validator was correct; its coverage was not.

rejectMalformedErrorPattern becomes rejectMalformedProfilePatterns and now
compiles both keys, reporting failures from either in one message.

Verified by me, not taken on the worker's word: the actual merge of this
branch into main builds green (1280 tests), and I reran the mutation myself
— narrowing the loop back to errorPattern turns exactly the two new tests
red, one of them with 'Expected IllegalStateException to be thrown, but
nothing was thrown', which is the defect stated out loud.
2026-09-04 10:09:15 +07:00
Dai Ha f04e934b94 fleetd #273: validate exhaustedPattern regex at load, like errorPattern
CI / build (pull_request) Successful in 1m36s
CI / contract (pull_request) Successful in 2m11s
FleetConfig.rejectMalformedErrorPattern only compiled errorPattern eagerly
at config load. exhaustedPattern was compiled unguarded in Fleetd.main,
so profiles.<name>.exhaustedPattern: "[" passed load() and then crashed
the whole daemon at boot with a raw PatternSyntaxException naming neither
the profile nor the key.

Rename the validator to rejectMalformedProfilePatterns and extend it to
also compile every non-blank exhaustedPattern, reporting
profiles.<name>.exhaustedPattern ("<value>"): <message> in the same style
as errorPattern. Both keys are collected and reported together from a
single load. Fleetd.java's compile site is left as-is per scope — it is
now safe because load already rejects a bad value.

Added tests covering: a bad exhaustedPattern is refused; a bad pattern in
each key is reported together in one message; valid patterns still load;
a blank/absent exhaustedPattern is ignored.
2026-09-04 10:05:36 +07:00
Dai Ha 3fae35c357 #276: say whose environment the "allowed N of M" line counted
CI / contract (push) Successful in 55s
CI / build (push) Successful in 1m42s
A gap in my own #269 fix. That ticket stopped four sites claiming things
about the member's environment that fleetd cannot see when memberHerdrSocket
is configured, and gave the WARN in logCredentialGap a guard. The INFO line
called two lines earlier never got one:

    logAllowListCoverage(allowed);      // no guard
    logCredentialGap(creds, allowed);   // guarded since #269

Read plainly, "member credentials: allowed 7 of 39" is a statement about the
member's credentials. Under memberHerdrSocket the pane is routed to a second
herdr whose environment fleetd has no channel to inspect, so those counts
come from fleetd's own process instead. Same overclaim #269 existed to
remove, in the line next door.

The method's javadoc does carry the caveat, by cross-reference to another
field's javadoc. That does not help the operator reading fleetd.out.

The counts stay useful, so this is not a WARN and not a refusal — only the
claim is narrowed. The unguarded path keeps its exact original wording, so
the existing assertion on "member credentials: allowed 1 of 3" still holds.

The new test pins the pair together so a later edit cannot fix one line and
leave the other. Mutation-proved: with the guard removed it fails printing
the old line verbatim.

Also worth recording: no test covered #269's own guard — that WARN wording
shipped unverified, and still has no coverage.
2026-09-04 10:01:23 +07:00
9 changed files with 381 additions and 23 deletions
+13
View File
@@ -39,6 +39,19 @@ a worker made all 59 of its edits in the primary's tree and never noticed.
test "$(git rev-parse --show-toplevel)" = "$PWD" || cd "$(git rev-parse --show-toplevel)"
```
**Never run `git stash` (or `git stash pop`/`apply`/`drop`).** Your worktree is isolated, but the
stash is **not**: `refs/stash` is one stack shared by the primary's checkout and every other
worker's worktree of this repo. Measured on 2026-09-04 — `git stash list` from a worker's worktree
and from the primary's tree returned byte-identical output. So a `git stash` you run can be popped
into someone else's tree, and a `git stash pop` you run can drop **another worker's** uncommitted
edits on top of yours. This has already happened here: two workers were running in parallel and one
of them had its in-progress edit silently overwritten by the other's stash.
The branch is your isolation, so use it instead. To set work aside, commit it on your own branch
(`git commit -m "wip: ..."`) and carry on; to try something and back out, use
`git diff > /tmp/<your-branch>.patch` then `git checkout -- <file>`. Both stay inside your worktree.
If you find a stash entry you did not create, leave it alone and say so in your report.
## 2. Implement
- Implement exactly the scope the lead named. Keep the diff focused; note anything out of scope
+76
View File
@@ -0,0 +1,76 @@
# Authorization entry-point audit
Scope reviewed: every route registered in `FleetApp.build`, every tool handler
registered in `FleetMcp`, and their shared `CallerResolver` and `Authz` gate.
## Result
No authorization-action mismatch was found. The only handler whose operation
changes with its arguments is `fleet_poll`. It selects `DRAIN` when `target` is
present and `READ` when it is absent, before it calls either service branch.
`Fleetd` creates one `CallerResolver` and passes that same instance to both
`FleetMcp` and `FleetApp` (`Fleetd.java:615-650`, `720-722`). REST resolves it
in the Javalin pre-handler. MCP resolves it in the transport context extractor.
## REST routes
| Route | Gate and choice location | Branch/target review | Verdict |
|---|---|---|---|
| `GET /healthz` | None | Liveness probe only; intentionally open. | ok |
| `GET /metrics` | `METRICS`, start of `metrics` | No branch or target. | ok |
| `GET /sessions` | `READ`, start of `sessions` | No caller-selected target. | ok |
| `GET /agents` | `READ`, start of `agents` | No caller-selected target. | ok |
| `GET /members` | `READ`, start of `listMembers` | No caller-selected target. | ok |
| `GET /profiles` | `READ`, start of `profiles` | No caller-selected target. | ok |
| `GET /member-credentials` | `READ`, start of `memberCredentials` | No caller-selected target. The view exposes policy names and counts, not values. | ok |
| `POST /members` | `SPAWN`, before query/body parsing in `spawnMember` | Arguments select role, profile, cwd, and worktree. They do not select a different authority type. | ok |
| `DELETE /members/{paneId}` | `STOP`, after reading `paneId` in `stopMember` | Caller can select another pane, but only a primary has `STOP`. | ok |
| `POST /sessions/{id}/message` | `SEND`, after reading path `id` in `sendMessage` | Normal send, async send, and `turnId` answer all deliver a turn/message. `turnId` does not widen the roles allowed to send. | ok |
| `POST /sessions/{id}/reply` | `REPLY`, after reading path `id` in `replyMessage` | Caller can name a target, and `Authz` requires it to equal the connection-resolved terminal. | ok |
| `GET /sessions/{id}/replies` | `DRAIN`, after reading path `id` in `drainReplies` | Removes inbox entries; only a primary has `DRAIN`. | ok |
| `POST /sessions/{id}/ask` | `ASK`, after reading path `id` in `askMessage` | Caller can name a target, and `Authz` requires it to equal the connection-resolved terminal. | ok |
| `GET /sessions/{id}/status` | `READ`, after reading path `id` in `sessionStatus` | Any authenticated role may observe any session. This matches the `READ` policy, which intentionally does not use target ownership. | ok |
| `GET /tasks/{ticket}` | `READ`, start of `taskStatus` | Ticket polling is read-only; no service branch changes the action. | ok |
## MCP tools
| Tool | Gate and choice location | Branch/target review | Verdict |
|---|---|---|---|
| `fleet_send` | `SEND`, at the start of `sendHandler` | `coordId`, `turnId`, synchronous, and async forms all deliver a message or a turn answer. `sessionId` is read before the branch for audit target only. | ok |
| `fleet_reply` | `REPLY`, after deriving the connection terminal in `replyHandler` | No target argument exists. The service always receives the caller's own terminal. | ok |
| `fleet_ask` | `ASK`, after deriving the connection terminal in `askHandler` | No target argument exists. The service always receives the caller's own terminal. | ok |
| `fleet_status` | `READ`, start of `statusHandler` | Any authenticated role may query any session. This is the same intentional `READ` policy as the REST route. | ok |
| `fleet_poll` with `ticket` | `READ`, `pollAction(target)` before dispatch in `pollHandler` | Reads a task only. | ok |
| `fleet_poll` with `target` | `DRAIN`, `pollAction(target)` before dispatch in `pollHandler` | Drains and removes a target inbox. Only a primary has `DRAIN`. | ok |
| `fleet_ack` | `DRAIN`, start of `ackHandler` | Removes one target inbox entry. Only a primary has `DRAIN`. | ok |
| `fleet_spawn` | `SPAWN`, start of `spawnHandler` | Arguments select member configuration only. | ok |
| `fleet_list` | `READ`, start of `listHandler` | No caller-selected target. | ok |
| `fleet_stop` | `STOP`, after reading `paneId` in `stopHandler` | Caller can select another pane, but only a primary has `STOP`. | ok |
| `fleet_profiles` | `READ`, start of `profilesHandler` | No caller-selected target. | ok |
| `fleet_whoami` | `READ`, start of `whoamiHandler` | Reports the connection-resolved caller, not an argument. | ok |
## Surface parity
The matching route/tool pairs use the same action:
| Operation | REST | MCP | Result |
|---|---|---|---|
| spawn | `SPAWN` | `SPAWN` | match |
| stop | `STOP` | `STOP` | match |
| send and ask answer | `SEND` | `SEND` | match |
| reply | `REPLY` | `REPLY` | match |
| ask | `ASK` | `ASK` | match |
| drain replies | `DRAIN` | `DRAIN` for `fleet_poll{target}` and `fleet_ack` | match |
| status | `READ` | `READ` | match |
| task/ticket polling | `READ` | `READ` | match |
| list/profiles | `READ` | `READ` | match |
## Test risk
`FleetMcpAuthzTest` has a direct regression test for the argument-dependent
`fleet_poll` action. It tests the shared role table for the other tools, but it
does not pin each handler's chosen action. This is not a current defect because
the reviewed handlers choose the matching action. A future change that adds an
argument-dependent operation should add a handler-level action-selection test,
like `pollingByTargetIsADrainAndPollingByTicketIsARead`.
@@ -593,7 +593,11 @@ public final class Fleetd {
if (detail.agentSessionId() != null) {
reason += " agentSessionId=" + detail.agentSessionId();
}
messages.abandon(detail.terminalId(), reason);
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
// (see MessageService.abandon's javadoc for why those two must differ).
messages.abandon(detail.terminalId(), reason, true);
replyInbox.release(detail.terminalId());
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
});
@@ -1506,7 +1506,7 @@ public record FleetConfig(
rejectDuplicateMemberSlots(yaml);
rejectNegativeMaxLoad(yaml);
rejectAutoCompactWindowOutOfRange(yaml);
rejectMalformedErrorPattern(yaml);
rejectMalformedProfilePatterns(yaml);
rejectUnknownKind(yaml);
rejectUnknownAuthMode(yaml);
rejectUnknownPlacement(yaml);
@@ -1856,20 +1856,26 @@ public record FleetConfig(
}
/**
* Reject a profile whose {@code errorPattern} (fleetd #201 Unit 5) is not a valid Java regex,
* naming the profile, the key, and the parser's own message.
* Reject a profile whose {@code errorPattern} (fleetd #201 Unit 5) or {@code exhaustedPattern}
* (CB-578 stage A) is not a valid Java regex, naming the profile, the key, and the parser's own
* message.
*
* <p>Unset/{@code null} means "use {@code CompletionResolver}'s built-in {@code (?i)\bAPI
* Error\s*:} compatibility pattern" and passes silently. A profile that DOES set the key gets it
* compiled once at daemon startup ({@code Fleetd.main}, mirroring {@code exhaustedPattern}) — an
* uncaught {@link java.util.regex.PatternSyntaxException} there crashes startup without naming
* which profile or key is at fault. Validate eagerly here instead, at config load, the same
* "fail loud at load, not lazily later" reasoning as {@link #rejectAutoCompactWindowOutOfRange}.
* <p>Unset/{@code null} means, for {@code errorPattern}, "use {@code CompletionResolver}'s
* built-in {@code (?i)\bAPI Error\s*:} compatibility pattern", and for {@code exhaustedPattern},
* "opt out of that classification" — either way it passes silently. A profile that DOES set
* either key gets it compiled once at daemon startup ({@code Fleetd.main}) — an uncaught
* {@link java.util.regex.PatternSyntaxException} there crashes startup without naming which
* profile or key is at fault (fleetd #273: this happened for {@code exhaustedPattern}, which had
* no validator here even though its sibling {@code errorPattern} did). Validate eagerly here
* instead, at config load, the same "fail loud at load, not lazily later" reasoning as
* {@link #rejectAutoCompactWindowOutOfRange}. Both keys are checked from a single load, and any
* failures from either are collected together into one message.
*
* @param yaml the raw config text
* @throws IllegalStateException when any profile's {@code errorPattern} fails to compile
* @throws IllegalStateException when any profile's {@code errorPattern} or
* {@code exhaustedPattern} fails to compile
*/
static void rejectMalformedErrorPattern(String yaml) {
static void rejectMalformedProfilePatterns(String yaml) {
Map<?, ?> raw;
try {
raw = YAML.readValue(yaml, Map.class);
@@ -1884,18 +1890,20 @@ public record FleetConfig(
if (!(e.getValue() instanceof Map<?, ?> p)) {
continue;
}
if (!(p.get("errorPattern") instanceof String pattern) || pattern.isBlank()) {
continue;
}
try {
Pattern.compile(pattern);
} catch (PatternSyntaxException ex) {
bad.add("profiles." + e.getKey() + ".errorPattern (\"" + pattern + "\"): " + ex.getMessage());
for (String key : List.of("errorPattern", "exhaustedPattern")) {
if (!(p.get(key) instanceof String pattern) || pattern.isBlank()) {
continue;
}
try {
Pattern.compile(pattern);
} catch (PatternSyntaxException ex) {
bad.add("profiles." + e.getKey() + "." + key + " (\"" + pattern + "\"): " + ex.getMessage());
}
}
}
bad.sort(String::compareTo);
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: malformed errorPattern — "
throw new IllegalStateException("refusing to start: malformed pattern — "
+ String.join("; ", bad));
}
}
@@ -1497,10 +1497,27 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* survive {@code allowed} (including the {@code LC_*} prefix rule). Neither number is a constant:
* both come from the actual derived set and the actual environment this spawn sees. Never logs a
* variable NAME or VALUE — only the counts.
*
* <p><strong>Under {@code memberHerdrSocket} the counts describe fleetd's own process, not the
* member's</strong> (fleetd #269 follow-up), so the message says so rather than leaving the
* reader to infer it from this javadoc, which the operator reading the log never sees.
*/
private void logAllowListCoverage(Set<String> allowed) {
Set<String> hostNames = hostEnvNames.get();
long kept = hostNames.stream().filter(name -> MemberEnvAllowList.keeps(allowed, name)).count();
if (memberHerdrSocketConfigured()) {
// fleetd #269 covered the sibling line below (logCredentialGap) and stopped there.
// This line has the same problem: read plainly, "allowed 7 of 39" is a statement about
// the member's pane, and under memberHerdrSocket it is not -- the pane is routed to a
// second herdr whose environment fleetd cannot inspect. The counts stay useful, so
// this is not a WARN and not a refusal; only the claim is narrowed to what is true.
log.info("member credentials: allowed {} of {} names in fleetd's OWN environment — "
+ "memberHerdrSocket is configured, so member panes are routed to a "
+ "second herdr whose environment fleetd has no channel to inspect. "
+ "These counts describe fleetd's process, NOT the member pane's.",
kept, hostNames.size());
return;
}
log.info("member credentials: allowed {} of {}", kept, hostNames.size());
}
@@ -360,6 +360,18 @@ public final class MessageService {
* {@link #ask} clears the ticket's question and returns it to {@code PENDING}, but {@link #send}
* already closed the forward waiter the instant the question surfaced, so the target has
* neither an accepted nor a queued delivery left to show for it.
*
* <p><strong>Deliberately still {@code question == null} only (fleetd #275).</strong> This
* method must not also report a still-{@link Phase#ASKING} task as orphaned: the worker may
* genuinely be waiting on a live primary that is about to (or already mid-{@link #answer})
* answer it, and {@link dev.ltms.fleet.health.FleetHealthMonitor} would classify that as
* {@code DELEGATION_ORPHANED} on nothing more than an active, healthy conversation. {@link
* #abandon(String, String, boolean)}'s {@code sweepAsking} path fixes the actual reachable gap
* (a target torn down for good while genuinely {@code ASKING}) at the point of teardown itself,
* by completing the task's future right there — so by the time this method would ever see it,
* {@code task.future.isDone()} is already {@code true} and it is excluded regardless of this
* guard. Widening this check instead of that one would trade a real fix for false positives on
* every ordinary in-flight question.
*/
public boolean hasOrphanedDelegation(String target) {
if (target == null || hasAcceptedDelivery(target) || hasQueuedDelivery(target)) {
@@ -580,6 +592,39 @@ public final class MessageService {
* reply — see the note above)
*/
public boolean abandon(String target, String reason) {
return abandon(target, reason, false);
}
/**
* As {@link #abandon(String, String)}, with control over whether a task still paused in
* {@code fleet_ask} ({@link Phase#ASKING}) is swept too (fleetd #275).
*
* <p>{@code sweepAsking} must be {@code true} only when the caller has independent, certain
* knowledge that {@code target} can never resume its turn — today that is only
* {@code sessions.onRelease}'s teardown (an explicit {@code fleet_stop}, or the idle reaper):
* the worker's pane is being stopped right now, so whatever it was mid-{@code fleet_ask} about
* has no turn left to resume into. {@link dev.ltms.fleet.health.FleetHealthMonitor}'s
* health-classification call keeps passing {@code false} (via {@link #abandon(String, String)}):
* a GONE/NEVER_READY reading is the daemon's best guess from the live agent list, not a teardown
* it performed itself, and {@code abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer} documents
* why an active ask must survive that guess — the primary may already be mid-{@link #answer} for
* the very same turn, and completing it here first would preempt a real answer with a misleading
* failure.
*
* <p><strong>Without {@code sweepAsking} on the release path, a target torn down while
* genuinely {@code ASKING} was unrecoverable.</strong> {@link #resolveQuestion} had already
* closed the forward waiter the instant the question surfaced (so the {@code waiter} branch
* below finds nothing to fail), the {@code question == null} guard excluded the task from
* {@code matching} (so the loop below skipped it too), and the worker's own {@code fleet_ask}
* clears {@link Task#question} back to {@code null} only once it lapses (the reverse-rendezvous
* window — up to {@code FleetMcp.ASK_DEFAULT_TIMEOUT_MS} / {@code FleetApp.MAX_ASK_TIMEOUT_MS},
* 55–115s) — by which point the released session no longer appears in {@code sessions.roster()}
* for {@link dev.ltms.fleet.health.FleetHealthMonitor} to ever re-observe, so nothing was ever
* left to call {@link #abandon} on this target again. The ticket then sat in {@link #tasks}
* forever: not terminal, so {@link #pruneTerminalTickets} never dropped it, and
* {@code fleet_poll} reported it stuck at {@link Phase#PENDING} for good.
*/
public boolean abandon(String target, String reason, boolean sweepAsking) {
boolean hadStrandedReply = hasStrandedReply(target);
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
strandedReplies.remove(target);
@@ -590,7 +635,8 @@ public final class MessageService {
List<Task> matching = new ArrayList<>();
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null && !task.future.isDone()) {
if (target.equals(task.target) && (sweepAsking || task.question == null)
&& !task.future.isDone()) {
matching.add(task);
}
}
@@ -607,11 +653,20 @@ public final class MessageService {
for (Task task : matching) {
boolean isRecovery = task == recoveryTask && recovered != null;
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
String turnId = task.turnId;
if (task.future.complete(outcome)) {
if (outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
} else if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
if (turnId != null) {
// #275: whether this task was swept out of ASKING or was already answered and
// only waiting on its resumed turn's real reply (#137), nothing will ever
// complete this turnId now — drop it from this class's own bookkeeping AND the
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
// that is actually done, and a late answer() sees it as lapsed rather than
// resolving a question nothing is listening for any more.
asyncTasksByTurn.remove(turnId, task);
rendezvous.closeAsk(turnId);
}
} else if (isRecovery) {
// The recovered reply was already drained out of the inbox, but this task resolved
@@ -155,6 +155,90 @@ class FleetConfigTest {
assertTrue(e.getMessage().contains("errorPattern"), "the offending key is named: " + e.getMessage());
}
// ── fleetd #273: exhaustedPattern gets the same load-time validation as its sibling errorPattern ──
@Test
void aProfileWithAMalformedExhaustedPatternIsRejectedAtLoadNamingTheProfileAndKey(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("malformed-exhausted-pattern.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
exhaustedPattern: "["
""");
IllegalStateException e = assertThrows(IllegalStateException.class, () -> FleetConfig.load(f));
assertTrue(e.getMessage().contains("ltms-local"), "the offending profile is named: " + e.getMessage());
assertTrue(e.getMessage().contains("exhaustedPattern"), "the offending key is named: " + e.getMessage());
}
@Test
void aMalformedErrorPatternAndAMalformedExhaustedPatternAreBothReportedFromOneLoad(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("both-malformed.yaml");
Files.writeString(f, """
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
errorPattern: "(unterminated["
terra:
baseUrl: http://gx01.gw:8000
exhaustedPattern: "["
""");
IllegalStateException e = assertThrows(IllegalStateException.class, () -> FleetConfig.load(f));
assertTrue(e.getMessage().contains("sonnet"), "the errorPattern profile is named: " + e.getMessage());
assertTrue(e.getMessage().contains("errorPattern"), e.getMessage());
assertTrue(e.getMessage().contains("terra"), "the exhaustedPattern profile is named: " + e.getMessage());
assertTrue(e.getMessage().contains("exhaustedPattern"), e.getMessage());
}
@Test
void validErrorPatternAndExhaustedPatternBothLoadFine(@TempDir Path dir) throws Exception {
Path f = dir.resolve("both-valid.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
errorPattern: "credential outage"
exhaustedPattern: "usage limit has been reached"
""");
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
assertEquals("credential outage", w.errorPattern());
assertEquals("usage limit has been reached", w.exhaustedPattern());
}
@Test
void aBlankExhaustedPatternNormalizesToNullJustLikeUnset(@TempDir Path dir) throws Exception {
Path f = dir.resolve("blank-exhausted-pattern.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
exhaustedPattern: " "
""");
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
assertNull(w.exhaustedPattern());
assertFalse(w.hasExhaustedPattern());
}
@Test
void aProfileWithNoExhaustedPatternLoadsFine(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-exhausted-pattern.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
""");
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
assertNull(w.exhaustedPattern());
assertFalse(w.hasExhaustedPattern());
}
@Test
void withProfileCarriesErrorPatternThrough(@TempDir Path dir) throws Exception {
Path f = dir.resolve("with-profile-error-pattern.yaml");
@@ -250,6 +250,57 @@ class HerdrPeerLauncherAllowListWiringTest {
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* fleetd #269 follow-up: the same overclaim the WARN in {@code logCredentialGap} was fixed for,
* in the INFO line beside it. With {@code memberHerdrSocket} configured, member panes are routed
* to a second herdr whose environment fleetd has no channel to inspect, so the counts come from
* fleetd's OWN environment. The bare line "member credentials: allowed 1 of 3" reads as a fact
* about the member's pane, and there it is not one.
*
* <p>#269 reworded four sites and stopped at the sibling below; this pins the pair together so
* a future edit cannot fix one and leave the other. Real path: asserted after a real {@link
* HerdrPeerLauncher#spawn}, reading the log production actually emits.
*/
@Test
void theAllowedCountLineSaysWhoseEnvironmentItCountedWhenMemberHerdrSocketIsSet(@TempDir Path worktreeRoot)
throws IOException {
String group = currentUserGroup();
FakeHerdr herdr = new FakeHerdr();
Set<String> hostEnvNames = Set.of(INJECTED, "SOME_UNRELATED_NAME", "ANOTHER_UNRELATED_NAME");
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/bash-should-be-ignored",
() -> hostEnvNames,
() -> configWithMemberHerdrSocketRootAndGroup("/tmp/other-user-herdr.sock", "/bin/zsh",
worktreeRoot.toString(), group));
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
List<String> lines = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
String coverage = lines.stream()
.filter(l -> l.startsWith("member credentials: allowed "))
.findFirst()
.orElse(null);
assertNotNull(coverage, "the coverage line must still be logged — narrowing the claim must "
+ "not silently delete the line: " + lines);
assertTrue(coverage.contains("fleetd's OWN environment"),
"the line must say whose environment it counted: " + coverage);
assertTrue(coverage.contains("NOT the member pane's"),
"and must say plainly that it is not the member's: " + coverage);
// The counts themselves stay real — narrowing the claim must not turn them into constants.
assertTrue(coverage.startsWith("member credentials: allowed 1 of 3"),
"the real counts must survive the rewording: " + coverage);
}
/**
* Lead-review fix: on a NON-zsh shell no scrub ever runs (bash ignores {@code ZDOTDIR}), so the
* "allowed N of M" line — which describes what the scrub does — must not be printed there either.
@@ -828,6 +828,56 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
// --- fleetd #275: a target torn down FOR GOOD while genuinely ASKING must not orphan --------
//
// sessions.onRelease (fleet_stop, or the idle reaper) is the one abandon() caller that knows
// for certain the target can never resume: its pane is being stopped right now. Unlike the
// health-classification caller above (a GONE/NEVER_READY guess, not a teardown it performed),
// it must sweep an ASKING ticket right here — see MessageService.abandon(String, String,
// boolean)'s javadoc for the full reachability chain this closes: without this, the forward
// waiter is already closed by the time the question surfaces, the ASKING guard skips the task,
// and by the time the worker's own fleet_ask lapses (~55-115s later) the released session no
// longer appears in FleetHealthMonitor's roster for anything to ever sweep it again — leaving
// fleet_poll{ticket} stuck PENDING forever.
@Test
void abandonWithSweepAskingFailsATornDownTargetsAskingTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 300));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertTrue(messages.abandon(T, "the worker session was released before it replied", true),
"a released target's open ask can never resume, so it must fail right here");
MessageService.TaskView failed = awaitTicketPhase(ticket, MessageService.Phase.FAILED);
assertEquals("the worker session was released before it replied", failed.detail());
// The reverse-rendezvous ask is torn down too: the worker's still-blocked fleet_ask rides
// out its own timeout (nothing completed its answer future), and a late answer() for the
// same turnId must see it as lapsed rather than resolving a question nobody is waiting on.
assertEquals(MessageService.AskOutcome.TIMED_OUT, ask.get(5, TimeUnit.SECONDS).outcome());
assertEquals(MessageService.Outcome.STALE_TURN,
messages.answer(asking.turnId(), "config.yaml", 200).outcome());
}
@Test
void abandonWithoutSweepAskingBehavesLikeTheTwoArgOverload() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertFalse(messages.abandon(T, "agent target term_a not found", false),
"sweepAsking=false must match the plain abandon(target, reason) overload");
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase());
}
// --- #137: a fleet_ask round-trip must not orphan the ticket's own reply -------------------
//
// The primary's fleet_send{turnId} answer call is itself bounded (a real MCP call, capped well