Compare commits

...

20 Commits

Author SHA1 Message Date
Dai Ha e5cb51a90e #324: read task.turnId once in finishAsyncTask to stop an NPE from ask()'s unlocked forgetting
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 2m1s
answer() holds sessionLocks while finishAsyncTask reads the volatile Task.turnId twice — once to
check it is non-null, once as the ConcurrentHashMap.remove key. ask()'s own timeout path mutates
the same field with no lock, via clearAsyncQuestion(turnId, true). volatile makes each read fresh
but not the pair atomic, so the field can go null between the two reads and remove(null, task)
throws NullPointerException on the lead's own answer() call, even though the reply already
completed on the line above.

Capture task.turnId into a local once and use that for both the check and the removal.

Added a package-private test seam (finishAsyncTaskRaceHook + forgetTurnForTest) so a test can force
the exact interleaving deterministically, by running the identical clearAsyncQuestion(turnId, true)
cleanup ask() uses, at the point between finishAsyncTask's former two reads. Both are inert (null)
in production.
2026-09-04 14:38:52 +07:00
Dai Ha fa1f49675b Merge #318: a delivery landing after release is refused, not parked in a map nobody reads
CI / contract (push) Successful in 1m0s
CI / build (push) Successful in 2m8s
2026-09-04 14:21:29 +07:00
Dai Ha 8426c3528f #316: pin the fail-toward-preserve rule on the late re-check, found by mutation
CI / contract (push) Successful in 1m22s
CI / build (push) Successful in 1m42s
2026-09-04 14:20:36 +07:00
Dai Ha 65f98ba910 Merge #316: the dirty check that authorises the worktree removal is taken after the worker stops 2026-09-04 14:17:13 +07:00
Dai Ha d05205d1eb Add a hunter skill: a sweep and a diff review are different jobs with different output contracts
CI / contract (push) Successful in 47s
CI / build (push) Failing after 1m49s
2026-09-04 14:13:08 +07:00
Dai Ha 667254df47 #316: re-check worktree dirtiness after the pane stops, before removing it
CI / contract (pull_request) Successful in 39s
CI / build (pull_request) Successful in 1m46s
SessionManager.releaseRemoved read hasUncommitted() once, while the worker
could still write, then used that stale boolean after launcher.stop() to
authorise `git worktree remove --force`. The same stale read also gated
trySnapshot, so a worker that wrote between the read and the stop lost its
work with neither a preserve nor a snapshot.

Add a second, best-effort hasUncommitted read immediately before the
removal, taken only on the path that is actually about to delete something
(never on a release that already decided to preserve, and never for
SHUTDOWN, which preserves unconditionally). If the tree is now dirty,
preserve it and attempt a fresh snapshot, since the original snapshot never
ran when the pre-stop read said clean. A failing re-check also preserves,
matching the existing CB-581 fail-safe rule.
2026-09-04 14:12:49 +07:00
Dai Ha 2926cd1784 #318: release() no longer strands a delivery that lands while it is running
CI / build (pull_request) Failing after 1m59s
CI / contract (pull_request) Successful in 2m14s
AmqpReplyInbox.release() used held.remove(target) then iterated the old
map. A delivery landing on the consumer work-pool thread after the
remove (basicCancel does not flush one already handed to that pool) hit
deliverCallback's computeIfAbsent, found the key gone, and created a
brand-new map release() never looks at again — delivered-but-unacked
forever, never requeued, never redelivered (#298 only closed the
"already in held when release runs" case).

Fix: release() swaps in a RELEASED tombstone via held.compute(...)
instead of held.remove(...). ConcurrentHashMap serializes compute/
computeIfAbsent calls for the same key against each other, so whichever
of release() and a concurrent deliverCallback runs first is fully
visible to the other — no gap. deliverCallback checks for the
tombstone and nacks-with-requeue instead of recreating a map; peek/ack
treat it as empty; own() clears a stale tombstone so a target is never
poisoned if its id is ever reused (the issue's own text says id reuse
doesn't happen, but the tombstone would otherwise sit in `held` forever
either way).

New test AmqpReplyInboxReleaseRaceTest forces the actual interleaving
with a latch (blocks release() inside its nack loop, which is only
reachable after the tombstone swap, then fires a concurrent delivery)
rather than a sequential call — a sequential test would not have caught
this, since #298's own contract test forces settlement before release()
runs. Mutation-tested: reverting the fix makes this test fail with
"expected: <2> but was: <1>" (m1 never nacked); restored after
confirming that failure.
2026-09-04 14:10:21 +07:00
Dai Ha c801851c66 Correct the isLoopback javadoc: after #305 narrowing this range refuses a caller, it does not promote one
CI / contract (push) Successful in 1m16s
CI / build (push) Successful in 1m47s
2026-09-04 14:10:10 +07:00
Dai Ha de70aa38f1 Merge #317: an unresolved caller is refused, never promoted to primary
CI / contract (push) Successful in 1m1s
CI / build (push) Successful in 1m56s
2026-09-04 14:03:35 +07:00
Dai Ha 53a533afb4 #317: refuse an unresolved caller instead of promoting it to primary
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 1m52s
ConnectionIdentity.resolve() called pids.pidForLocalPort(remotePort),
which returns -1 both on a real failure and (silently, no log line)
when lsof just finds no matching process. terminalForPid(-1) then
matches no pane, so CallerResolver's loopback-trust fallback could not
tell that caller apart from a genuine primary and handed it
Principal.primary(...) — granting SPAWN, STOP, SEND and DRAIN to a
worker whose PID lookup failed. This is the escalation PaneLocator's
own javadoc already names; CB-161's ancestry walk only helps once a
candidate pid exists, and a failed lookup has none.

Fix: ConnectionIdentity.Caller gets a resolved() predicate (pid > 0),
centralised next to the -1 sentinel it tests for the same reason
isLoopback() is centralised (fleetd #305: two independent copies of
one rule already drifted once). CallerResolver's loopback-trust
fallback now requires c.resolved() before granting PRIMARY; an
unresolved caller gets Principal.anonymous() — the same already-tested
"authenticated as nothing" outcome used everywhere else in that
method, so the refusal is a clean, named, unsurprising result rather
than something that looks like a bug.

Also logs the previously-silent "lsof ran clean, found no match" case
in LsofPeerPidLookup at DEBUG, since that (not a slow lsof — the
waitFor result was already discarded) is the likelier real trigger.

loopbackTrustTreatsANonWorkerLoopbackCallerAsThePrimary is untouched
and still green: a real pid that owns no pane (the actual primary) is
still resolved() and still PRIMARY. Token mode is unaffected — it
never consults c.pid() at all.

Mutation-tested: reverting only the CallerResolver.java guard
reproduces the escalation exactly (aFailedPeerPidLookupIsRefusedNotPromotedToPrimary
fails with "expected: <ANONYMOUS> but was: <PRIMARY>").
2026-09-04 13:58:14 +07:00
Dai Ha 77ad88631b Merge #315: the fixed placement policy honours the retry loop's unreachable set
CI / build (push) Successful in 1m22s
CI / contract (push) Successful in 1m55s
2026-09-04 13:57:24 +07:00
Dai Ha d88017807b #315: fix self-contradicting javadoc left by the previous commit
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 2m31s
FixedPlacementPolicy's class javadoc still opened with "This ignores caps
and reachability" after the previous commit added reachability as the
fourth carve-out that is explicitly NOT ignored — caught by a shape-check
survey run against this same file as part of #315's own request ("look in
placement/ ... for the same shape: a caller/comment that documents an
expectation ... where an implementation does not meet it"). Reworded the
opening sentence: fixed still ignores caps (maxLoad) by design, but
reachability is now a narrower, per-call retry exclusion, not an ignored
concern.
2026-09-04 13:55:28 +07:00
Dai Ha 2159a5a94a #315: FixedPlacementPolicy now honors the retry loop's unreachable set
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Successful in 2m33s
CompositePeerLauncher.spawn retries a failed candidate on the next one and
rebuilds PlacementContext "so the policy excludes this profile" (its own
comment), but FixedPlacementPolicy.select never read ctx.unreachable(). Under
the default `fixed` placement policy (used when `placement` is unset or set
to `fixed`), every retry re-picked the same dead default and a second,
healthy, configured profile was never tried. This also covers the wiring-bug
branch (a candidate profile with no owning adapter), which hit the exact same
symptom for the same reason.

Not live on this fleet: fleetd.yaml sets placement: weighted, which already
consults ctx.unreachable() via PlacementPolicyUtil.available(). This is live
only for a deployment that leaves placement unset or sets it to fixed.

Fix is in FixedPlacementPolicy: consult ctx.unreachable() in the same two
places it already consults quarantined/coolingOff (the default check and the
fallback walk over candidates()), and add a fourth reason to the "no
candidate remains" exception. Considered fixing this in
CompositePeerLauncher's retry loop instead (break when select() returns an
already-unreachable profile), but that only fails faster on the same dead
profile — it cannot make the loop advance to a different candidate, because
only the policy decides which candidate is next. The defect is that one
policy implementation does not honor the loop's stated contract, so the fix
belongs in that policy, matching how weighted/round-robin already behave.

Also fixed: the "no reachable worker profile" exception message said
"trying N candidate(s)" where N was unreachable.size(), a count of DISTINCT
profiles (a HashSet dedupes a profile added twice), under wording that reads
as a count of attempts. Reworded to "N distinct candidate(s)" so the count
matches what is measured and the profile list that follows it.

Tests: two new failover tests next to the three existing ones in
CompositePeerLauncherTest (which all use PlacementPolicies.weighted(), which
is why this had no coverage) — one pinned to PlacementPolicies.fixed() for
the unreachable-default case, one for the wiring-bug (no adapter) case.
Mutation-proofed: reverted FixedPlacementPolicy.java, both new tests failed
with the exact bug ("no reachable worker profile available after trying 1
distinct candidate(s): a" / "...c"), then restored the fix.
2026-09-04 13:52:29 +07:00
Dai Ha b9c2cf69f4 Merge #307: a worker's real reply after an ask timeout completes its ticket instead of stranding
CI / build (push) Successful in 2m5s
CI / contract (push) Successful in 2m19s
2026-09-04 13:31:43 +07:00
Dai Ha 8beae50fe7 Merge #308: refuse spawns once the shutdown drain has started, and sweep stragglers
CI / contract (push) Successful in 1m49s
CI / build (push) Successful in 3m0s
2026-09-04 13:28:09 +07:00
Dai Ha b2a58cb966 Merge #309: clean up partial worktree state when git worktree add fails 2026-09-04 13:28:04 +07:00
Dai Ha b8b25cf74c #307: an ask() timeout no longer strands the worker's real reply
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m30s
MessageService.reply()'s async-recovery path (askAnsweredAsyncTasks)
required a live Task.turnId, but ask()'s own TimeoutException handler
calls clearAsyncQuestion(turnId, true) — deliberately forgetting turnId
so hasAsyncQuestion() stops reporting the target BUSY. That made a
worker's eventual real fleet_reply, after an unanswered fleet_ask, fall
through to the inbox: fleet_poll{ticket} stayed PENDING forever and was
later force-failed with the false reason "session released before it
replied".

Fix: a new Task.askTimedOut marker is set (markAskTimedOut) right
before the turnId is forgotten, and askAnsweredAsyncTasks accepts it in
place of a live turnId. The marker never touches asyncTasksByTurn, so
the BUSY-release behaviour (invariant 1) is untouched. The existing
ambiguity guard (candidates.size() > 1 -> inbox, never guess) still
applies unchanged, but is now genuinely reachable rather than pure
defence in depth, since an ask timeout frees its target for a fresh,
independent delegation — the affected javadocs are updated to say so.

Tests: MessageServiceTest.aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket
(positive, mutation-proven) and
.twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess
(negative/ambiguity). FleetMcpTest's
unansweredAsyncAskReturnsTheTicketToPending was renamed and its final
assertion updated — it had pinned the old (buggy) inbox-stranding
behaviour as expected.
2026-09-04 13:23:01 +07:00
Dai Ha f159ca7d27 #310: log when a reap is skipped because the record changed
CI / contract (push) Successful in 1m7s
CI / build (push) Successful in 1m46s
The compare-and-release declines silently. This race is unobservable by
construction, so a reaper that quietly stops reaping is the hardest kind of
behaviour to diagnose later. One debug line names the pane and the likely
cause.
2026-09-04 13:22:55 +07:00
Dai Ha 83f2aea60f #308: refuse a spawn once the shutdown drain has started, and sweep stragglers
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 1m36s
drainAll iterated a one-shot registry snapshot with nothing to refuse a new
fleet_spawn while the drain was still running (mcp.close() only runs 8 calls
after sessions.close() in the shutdown hook). A session registered in that
window was never visited by the drain loop: its pane kept running and its
worktree was never preserved, with the in-memory registry gone at exit.

Fix, both mechanisms as the issue asked for (neither alone is complete):

- SessionManager.acquire now checks a `draining` flag, flipped true at the
  very start of drainAll before the registry snapshot is even taken, and
  throws the new ShuttingDownException (invariant 3: fail loudly, say why).
  FleetMcp.spawn and FleetApp.spawnMember surface it as a clean error/503
  rather than an uncaught RuntimeException.
- The flag alone cannot close the whole race: a caller already past the
  check can still be mid-launcher.spawn() (a real herdr round trip) when
  drainAll snapshots the registry. drainAll now re-reads the registry once
  its main pass finishes and drains whatever straggler landed there too,
  bounded by the SAME whole-drain deadline (invariant 1: timeoutNanos stays
  a budget for the whole drain, never extended for a straggler).
- ReleaseCause.SHUTDOWN still preserves worktrees for both the initial pass
  and the sweep (invariant 2, unchanged release() path).

Tests (SessionManagerTest): a guard test proving acquire() throws once
drainAll has started, and a race test using a launcher double that blocks
the second spawn() and the first stop() call to force, deterministically,
the exact interleaving where a spawn passes the guard before drainAll flips
it and only registers after the initial snapshot — proving the post-loop
sweep catches it.

Shape check (SessionManager.java only, not fixed): reapIdle has the same
shape — a decision made from a roster() snapshot, then acted on via
release(s.paneId()) with no re-check of the session's current state.
2026-09-04 13:21:53 +07:00
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
21 changed files with 1447 additions and 84 deletions
+102
View File
@@ -0,0 +1,102 @@
---
name: hunter
description: Defect-hunt procedure for a fleetd worker — sweep an assigned package for real bugs and report several ranked findings without fixing anything. Load this when the lead asks you to hunt or audit a scope rather than review one diff. Do NOT load `reviewer` for this; the two want different output.
---
# Hunter worker — procedure
The turn contract (one `fleet_reply`, `fleet_ask` for the lead's decisions, honest reporting,
never merge) is in **`CLAUDE.md` → Bridge communication → Worker** and already applies.
**This skill is not `reviewer`.** `reviewer` judges one diff and reports the *single* most
important issue in about 90 words. A hunt sweeps a whole package and reports *several* findings
in a long structured form. Loading both gives you two contradictory output contracts, and the
usual result is a worker that writes a good report into its terminal and ends the turn without
sending it. Load exactly one.
## 0. Read this before you read code: how the report gets home
Your terminal reaches nobody. The lead sees **only** the text inside your `fleet_reply` call.
A long report is exactly the case where this goes wrong, so plan for it:
- **Write the report into the `fleet_reply` argument itself.** Do not compose it in your terminal
and then summarise it into the call.
- If the report is long, **send it anyway** — one `fleet_reply` with everything.
- If you end the turn without replying, the bridge scrapes your pane instead. That scrape carries
at most the last 4000 characters, and on a hunt it usually captures the tail of the lead's own
brief rather than your findings. The lead then has nothing and has to ask you again.
## 1. Change nothing
A hunt is read-only. Do not edit a production file, do not "quickly fix" what you find, and do
not run a formatter. You may run the build and tests to *check* a claim, and you should say so
when you did.
## 2. Read the whole scope first
Read every file in the assigned package before you judge any of it. A defect that a caller
elsewhere in the same package makes unreachable is not a defect, and you cannot know that from
one file.
Stay inside the scope. If a defect there depends on a class outside it, read that class to
confirm — but the defect itself must live in the scope you were given.
## 3. The bar — this matters more than the count
**Name the path into the bad state.** Say which caller, in which state, reaches it. A defect on
paper is not a reachable defect. If you cannot name that path, keep the finding but mark it
`unproven` and say exactly what you could not check. Do not drop it, and do not dress it up.
**Say which direction the harm goes.** Data loss, privilege escalation and silent wrong answers
are worth reporting even when the window is narrow. A finding whose worst outcome is a worse log
line is not worth a block.
Two workers once ran the same scope: the one that applied the direction-of-harm filter found ten
real defects, the one that did not found none. Fewer findings the lead can act on beat many the
lead has to triage.
## 4. Shapes that have produced real merged fixes here
Read for these first:
1. **A one-way gate.** A guard added after an incident closes only the direction that incident
came from. Do not only ask what closes the gate — ask **which states still open it**.
2. **A value read once, then used later to authorise something destructive**, after something
else has had a chance to change it.
3. **A failure downgraded to a value that looks like a legitimate result** — `-1`, `null`, an
empty list, `false` — which a caller then trusts.
4. **A lock held for one half of a read-modify-write and not the other**, or two collections
updated under different locks.
5. **A comment or javadoc stating an invariant the code no longer keeps.** Comments are
load-bearing in this repo; a stale one has already caused a bug.
## 5. What you cannot check, and must not claim you did
- `fleetd/fleetd.yaml` is gitignored and **absent from your worktree**. You cannot read it. If a
finding depends on live configuration, name the key and say you could not check it.
- `.mcp.json`, `opencode.json` and `.autoenv` in your worktree are neutralised stubs, not the
repo's real files.
- The `wiki/` submodule pointer is months old. Do not cite it.
Reporting a fact you took from the lead's brief as something you measured yourself is a false
report, even when the fact is correct. Say where each fact came from.
## 6. The report — what goes in `fleet_reply`
One block per finding, most severe first:
```
FINDING N — <one line>
file:line
Path in: <which caller, in which state, reaches this>
Direction: <data loss | escalation | silent wrong answer | outage | ...>
Window/trigger: <when it actually happens>
Confidence: <confirmed by reading | unproven — say what you could not check>
Why nothing else catches it: <the guard or test you checked, and why it misses>
```
End with one line naming every file you read, so the lead knows the denominator.
**Nothing clears the bar?** Reply `NO FINDINGS`, name the files you read, and say what you ruled
out. A clean sweep is a valid result; an invented defect is worse than none.
+5
View File
@@ -10,6 +10,11 @@ never merge) is in **`CLAUDE.md` → Bridge communication → Worker** and alrea
skill is only the *review procedure*: how to work the scope, and the exact shape of what you
send back.
**Wrong skill for a sweep.** This one reviews *one* diff or scope and reports the *single* most
important issue. If the lead asked you to hunt or audit a whole package for several defects, load
`hunter` instead and ignore this file — the two want different output, and following both is how a
worker ends its turn with a good report that never gets sent.
## 1. Read the whole scope before you judge
The delegation names your scope — a file, a diff, a PR, a function. **Read all of it first.**
+7 -2
View File
@@ -200,8 +200,13 @@ must obey belongs in the charter, not here.
adapter, with a message naming the credential and the remaining seconds ("cooling off after
repeated backend errors") — distinct wording from a quarantine refusal, so don't conflate the
two when reading a spawn failure.
- **Skills available to delegate:** `implementer` (worktree → commit → push → own PR) and
`reviewer` (scoped review → one structured finding). Name one in every delegation.
- **Skills available to delegate:** `implementer` (worktree → commit → push → own PR),
`reviewer` (one diff → one structured finding) and `hunter` (sweep a package → several ranked
findings, change nothing). Name exactly one in every delegation. **`reviewer` and `hunter` are
not interchangeable** — `reviewer` caps the answer at one finding in about 90 words, so naming
it for a multi-finding sweep hands the worker two contradictory output contracts. That has
already cost three workers' turns: each wrote a good report to its terminal and ended the turn
with no `fleet_reply`, and the scrape returned the tail of the brief instead.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
`port-to-opencode` (make an OpenCode session a participant in this workspace) and
`fleets-status` (report every fleet that shares one LavinMQ instance).
@@ -236,7 +236,15 @@ public final class CallerResolver {
// loopback-trust: same-host callers that are not workers are the primary. A non-loopback
// caller is anonymous even here — and startup refuses that combination anyway
// (FleetConfig.validateAuthExposure), so this is defence in depth, not the control.
return isLoopback(remoteAddr) ? Principal.primary(c.pid()) : Principal.anonymous();
//
// fleetd #317: "not a worker" must not be conflated with "identity unresolved". The real
// primary is a real process — its pid resolves (c.resolved()), it just owns no herdr pane.
// A caller whose peer-PID lookup failed (LsofPeerPidLookup's -1 sentinel — on any failure,
// silently including "lsof found no match") has no such pid, and PaneLocator's own javadoc
// already names what happens if that case is handed the primary role: a worker→primary
// escalation. So an unresolved caller is refused (ANONYMOUS — the same clean, already-tested
// "authenticated as nothing" outcome used everywhere else in this method), never promoted.
return isLoopback(remoteAddr) && c.resolved() ? Principal.primary(c.pid()) : Principal.anonymous();
}
private boolean presentedTokenMatches(String authorizationHeader) {
@@ -36,6 +36,25 @@ public final class ConnectionIdentity {
* primary / an off-host client) and its {@code pid} (or {@code -1} if not resolvable).
*/
public record Caller(String terminal, long pid) {
/**
* Whether the OS peer-PID lookup actually succeeded — {@code false} means {@code pid} is
* the {@code -1} sentinel, not a real process id, so this caller's identity could not be
* established at all. That is a different fact from a real pid that simply owns no worker
* pane (the primary's own connection): the primary is {@code resolved()} and has a
* {@code null terminal}; an unresolvable caller is {@code !resolved()} and also has a
* {@code null terminal}. The two look identical through {@link #terminal} alone, which is
* exactly how fleetd #317 happened — a failed {@code lsof} lookup and a genuine primary both
* fell through to {@code Principal.primary(...)}.
*
* <p>Centralised here, next to the sentinel it tests, for the same reason
* {@link ConnectionIdentity#isLoopback} is centralised rather than left for each caller to
* reimplement: a raw {@code pid > 0} check duplicated at every call site is precisely the
* "one rule, two copies" shape that let #305 drift.
*/
public boolean resolved() {
return pid > 0;
}
}
/** Resolve the caller's terminal and PID from one peer-PID lookup. */
@@ -71,13 +90,27 @@ public final class ConnectionIdentity {
* {@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.
* <p><strong>What excluding an address costs, stated as it is today.</strong> This paragraph
* used to say that narrowing this range turned a worker into the lead, and that widening the
* check was what closed the hole. That was true only while there were <em>two</em> definitions
* that disagreed: {@code ConnectionIdentity} skipped the identity lookup for {@code 127.0.0.2}
* while {@code CallerResolver} read the same address as loopback and granted the primary role.
* #305 removed the second copy, and with one shared definition the old sentence no longer holds.
*
* <p>Measured on 2026-09-04 by narrowing this method back to exactly {@code 127.0.0.1} and
* running {@code CallerResolverTest} and {@code ConnectionIdentityTest}: a caller from
* {@code 127.0.0.2} then resolves to {@code ANONYMOUS}, not {@code PRIMARY} — for a worker
* ({@code aWorkerOnAnyLoopbackSourceAddressIsStillAWorkerNotThePrimary}) and for a non-worker
* ({@code aNonWorkerOnAnyLoopbackSourceAddressIsStillThePrimary}) alike. Excluding an address
* now <em>refuses</em> its caller; it does not promote one.
*
* <p>So keep the whole range, but for the plain reason: a genuine worker or primary that
* connects from {@code 127.0.0.2} must be identifiable at all, and narrowing this predicate
* locks it out. That is an outage, and an outage is the direction to fail in — which is exactly
* why the range must not be narrowed casually and also why doing so is no longer a security
* hole. This predicate still does not decide whether a caller is trusted; it decides whether the
* caller's identity is <em>resolved at all</em>. What makes an unresolved caller safe is
* {@link Caller#resolved()} (#317), not this method.
*/
public static boolean isLoopback(String addr) {
if (addr == null) {
@@ -22,6 +22,7 @@ import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.ShuttingDownException;
import dev.ltms.fleet.session.WorktreeRequest;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
@@ -962,6 +963,10 @@ public final class FleetMcp {
return text(json(memberView(member)));
} catch (GuardException e) {
return error("subscription boundary: " + e.getMessage());
} catch (ShuttingDownException e) {
// fleetd #308: the daemon's shutdown drain has already started — refuse loudly rather
// than register a session drainAll will never see again.
return error("shutting down: " + e.getMessage());
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct
// from "profile does not exist" below.
@@ -41,6 +41,14 @@ public final class LsofPeerPidLookup implements PeerPidLookup {
if (!p.waitFor(2, TimeUnit.SECONDS)) {
p.destroyForcibly();
}
if (found < 0) {
// fleetd #317: this is the silent path — lsof ran clean and simply reported no
// matching process (e.g. queried before the OS socket table settles). Previously
// this logged nothing at all, which is exactly why the escalation went unnoticed;
// the exception path below already logs. A caller now refused because of this is
// still refused (never promoted) — this line only makes the refusal diagnosable.
log.debug("lsof peer-pid lookup for port {} found no matching process", port);
}
return found;
} catch (Exception e) {
log.debug("lsof peer-pid lookup for port {} failed: {}", port, e.getMessage());
@@ -385,9 +385,12 @@ public final class CompositePeerLauncher implements PeerLauncher {
}
}
// unreachable.size() counts DISTINCT profiles, not attempts (a HashSet dedupes a profile
// added twice) — say "distinct" so the count matches the sentence and the profile list that
// follows, rather than reading as a count of attempts made (fleetd #315).
throw new PeerUnreachableException(
"no reachable worker profile available after trying " + unreachable.size()
+ " candidate(s): " + String.join(", ", unreachable));
+ " distinct candidate(s): " + String.join(", ", unreachable));
}
/**
@@ -22,6 +22,7 @@ import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicReference;
/**
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
@@ -93,8 +94,23 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private final Channel channel;
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
private final Object channelLock = new Object();
/** target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself. */
/**
* target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself.
*
* <p><strong>CB-318 tombstone.</strong> The value {@link #RELEASED} is a reserved sentinel: it
* marks a target whose {@link #release} has already run, so {@link #deliverCallback} can tell a
* delivery landing after release() apart from a fresh target it has never seen. See both methods'
* javadoc for why a plain {@code held.remove(target)} is not enough.
*/
private final ConcurrentHashMap<String, LinkedHashMap<String, Held>> held = new ConcurrentHashMap<>();
/**
* CB-318 sentinel stored in {@link #held} for a target whose {@link #release} has already run.
* Never mutated — every read site compares it by reference ({@code ==}) before touching it as a
* map, because it is a single object shared across every released target and calling a mutator on
* it would corrupt state for all of them.
*/
private static final LinkedHashMap<String, Held> RELEASED = new LinkedHashMap<>();
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
@@ -201,6 +217,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
String tag = channel.basicConsume(queue, false, deliverCallback(target), _ -> { });
consumerTags.put(target, tag);
// CB-318: drop a stale RELEASED tombstone from a prior ownership of this same target
// string, so a delivery under this fresh consumer is held normally instead of being
// nacked forever by deliverCallback's RELEASED check. Safe to do here, still under
// channelLock: no delivery for the consumer tag just registered above can reach
// deliverCallback before this basicConsume call returns.
held.remove(target, RELEASED);
log.debug("AMQP inbox owns queue {} for target {}", queue, target);
} catch (IOException e) {
throw new IllegalStateException("cannot own queue " + queue, e);
@@ -243,6 +265,32 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
* 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.
*
* <p><strong>CB-318: {@code held.remove(target)} alone leaves a second window open.</strong> The
* bullet above already explains why cancelling first does not save a tag from going stale — but
* that only accounts for a delivery landing before this method starts touching {@link #held}.
* {@code basicCancel} stops <em>new</em> dispatches; it does not flush one already handed to the
* consumer work pool. So a delivery can still land on that pool's thread and reach
* {@link #deliverCallback} at any point during, or after, this method's body — and a plain
* {@code held.remove(target)} does nothing to stop it: {@code deliverCallback}'s
* {@code computeIfAbsent} finds the key gone and happily creates a brand-new map under it, which
* this method — already past its {@code remove} — never looks at again. That entry then sits
* delivered-but-unacked on {@link #channel} until the whole inbox closes: never requeued, never
* redelivered, and {@link #peek} is never called again for a target nothing owns any more.
*
* <p>The fix is {@link #held}{@code .compute(target, ...)} instead of {@code remove}: it takes
* whatever was held (to nack, same as before) and, in the same atomic step, leaves the
* {@link #RELEASED} tombstone behind instead of an absent key. {@code computeIfAbsent} and
* {@code compute} calls for the same key are mutually exclusive in {@link ConcurrentHashMap} —
* whichever of this call and a concurrent {@code deliverCallback} runs first is fully visible to
* the other, with no gap between them. So a delivery that loses the race sees a real map here and
* gets nacked by the loop below, same as always; a delivery that wins the race (runs first) is
* itself nacked by that same loop, once it settles into {@code held}. A delivery that arrives once
* this method has stored {@link #RELEASED} finds it via {@code computeIfAbsent} and refuses itself
* — see {@link #deliverCallback}. Either way nothing is silently retained forever, satisfying the
* ticket's invariant against dropping a message. This closes the window rather than merely
* narrowing it — correctness does not depend on how much time elapses between the swap and this
* method returning.
*/
@Override
public void release(String target) {
@@ -255,8 +303,13 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
throw new IllegalStateException("cannot cancel consumer for " + target, e);
}
}
var perTarget = held.remove(target);
if (perTarget != null) {
AtomicReference<LinkedHashMap<String, Held>> previouslyHeld = new AtomicReference<>();
held.compute(target, (_, v) -> {
previouslyHeld.set(v);
return RELEASED;
});
var perTarget = previouslyHeld.get();
if (perTarget != null && perTarget != RELEASED) {
synchronized (perTarget) {
for (Held h : perTarget.values()) {
try {
@@ -326,7 +379,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
@Override
public List<InboxMessage> peek(String target) {
var perTarget = held.get(target);
if (perTarget == null) {
if (perTarget == null || perTarget == RELEASED) {
return List.of();
}
synchronized (perTarget) {
@@ -337,7 +390,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
@Override
public void ack(String target, String msgId) {
var perTarget = held.get(target);
if (perTarget == null) {
if (perTarget == null || perTarget == RELEASED) {
return;
}
Held h;
@@ -370,6 +423,20 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
String content = new String(delivery.getBody(), StandardCharsets.UTF_8);
var perTarget = held.computeIfAbsent(target, _ -> new LinkedHashMap<>());
if (perTarget == RELEASED) {
// CB-318: release() already ran for this target and left the RELEASED tombstone in
// held (see release()'s javadoc) — computeIfAbsent() is guaranteed to see it rather
// than recreate a fresh map, because ConcurrentHashMap serializes compute/
// computeIfAbsent calls for the same key against each other. Refuse the delivery
// instead of holding it somewhere release() will never look at again: requeue it, the
// same way release() nacks its own held entries, so a later owner (or a connection
// drop) can still recover it. This does not need channelLock across a broker round
// trip — basicNack, like the duplicate-ack case just below, does not wait for one.
synchronized (channelLock) {
channel.basicNack(tag, false, true);
}
return;
}
boolean duplicate;
synchronized (perTarget) {
if (perTarget.containsKey(msgId)) {
@@ -187,6 +187,24 @@ public final class MessageService {
private volatile Long completedNanos;
private volatile Reply question;
private volatile String turnId;
/**
* Set when this task's {@code fleet_ask} lapsed with no answer (fleetd #307):
* {@link #clearAsyncQuestion} then forgets {@link #turnId} (nulls it and drops the task from
* {@code asyncTasksByTurn}) so {@link #hasAsyncQuestion} stops reporting the target BUSY — a
* later {@code fleet_send} to it must be accepted, not refused. But the worker's turn is
* still genuinely live: it resumed on its own and will eventually call its real
* {@code fleet_reply}. Losing {@link #turnId} loses {@link #askAnsweredAsyncTasks}' only
* signal that such a reply belongs to this task, so that reply used to fall straight to the
* inbox and strand — {@code fleet_poll} stayed {@code PENDING} forever, later force-failed by
* {@link #abandon} with the misleading "session released before it replied". This flag is a
* second, independent signal that survives the forgetting: {@link #askAnsweredAsyncTasks}
* accepts it in place of a live {@link #turnId}, without ever re-adding the task to
* {@code asyncTasksByTurn} (so the BUSY release is untouched). Cleared implicitly once
* {@link #future} resolves — every match in {@link #askAnsweredAsyncTasks} already requires
* {@code !future.isDone()}, so a task that recovered (or was later failed by
* {@link #abandon}) can never match again regardless of this flag's value.
*/
private volatile boolean askTimedOut;
private Task(String ticket, String target, LongSupplier nowNanos) {
this.ticket = ticket;
@@ -395,18 +413,16 @@ public final class MessageService {
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
* are interactive and must never be queued.
*
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks}
* cannot actually return more than one entry today (see its own javadoc for why — in short,
* {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as
* long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from
* two other facts, not one this method enforces, so this branch stays in as defence in depth
* rather than being removed as dead code: if it ever weakens, returning whichever candidate a
* {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the
* <em>wrong</em> ticket — silently handing the lead something that reads like a correct answer to
* a delegation the worker never touched, which is worse than a failure because the lead acts on
* it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
* deterministically instead.
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks} can
* return more than one entry — a reachable state, not a hypothetical one (see its own javadoc:
* an {@code fleet_ask} that lapsed with no answer, fleetd #307, frees the target for a completely fresh
* delegation, which can itself go on to ask-and-lapse before the first worker's real reply
* arrives). Returning whichever candidate a {@code ConcurrentHashMap} iteration reaches first
* would let a genuine reply complete the <em>wrong</em> ticket — silently handing the lead
* something that reads like a correct answer to a delegation the worker never touched, which is
* worse than a failure because the lead acts on it. When more than one candidate exists, guessing
* is not safe: fall back to the inbox exactly 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
@@ -434,15 +450,18 @@ public final class MessageService {
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
}
// #137: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming a
// turn that {@link #answer} already gave up waiting on. answer()'s own bounded wait (the
// primary's fleet_send{turnId} call, capped well under a minute) can time out and close its
// waiter long before the worker — now actually resuming real work — finishes and replies. That
// reply used to have nowhere to land but the session inbox, leaving the async ticket's future
// unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's abandon() forced it
// FAILED with a misleading "session released before it replied" reason, even though the reply
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
// sees the real reply instead.
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
// - answer()'s own bounded wait (the primary's fleet_send{turnId} call, capped well under a
// minute) can time out and close its waiter long before the worker — now actually resuming
// real work — finishes and replies.
// - ask()'s own wait for the primary can time out first, with the worker resuming on its own
// and finishing unanswered.
// Either way that reply used to have nowhere to land but the session inbox, leaving the async
// ticket's future unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's
// abandon() forced it FAILED with a misleading "session released before it replied" reason,
// even though the reply had, in fact, arrived. Completing the matching ticket directly here
// means fleet_poll{ticket} sees the real reply instead.
List<Task> candidates = askAnsweredAsyncTasks(session);
if (candidates.size() == 1) {
Task orphan = candidates.get(0);
@@ -473,34 +492,42 @@ public final class MessageService {
}
/**
* Every still-open async task on {@code target} whose {@code fleet_ask} was already answered —
* its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer}
* — yet whose future is not resolved yet (#137). Empty if no such task exists, including the
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
* never asked has {@code turnId == null}, so it can never match here and only ever completes
* through the ordinary rendezvous fast path in {@link #reply}).
* Every still-open async task on {@code target} whose worker is genuinely expected to send a
* real {@code fleet_reply} next with nothing left registered to catch it: either its
* {@code fleet_ask} was already answered — {@link Task#turnId} is stamped but {@link
* Task#question} was cleared by {@link #answer} — or its {@code fleet_ask} lapsed unanswered and
* {@link Task#askTimedOut} marks that (fleetd #307; {@link Task#turnId} is {@code null} by then, forgotten
* so the target is not left BUSY — see {@link Task#askTimedOut}'s own javadoc). Either way the
* task's future is not resolved yet. Empty if no such task exists, including the common case
* where {@code target}'s worker never used {@code fleet_ask} at all (a task that was never asked
* has both {@code turnId == null} and {@code askTimedOut == false}, so it can never match here and
* only ever completes through the ordinary rendezvous fast path in {@link #reply}).
*
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@link #send}
* refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that
* check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} —
* not only while its question is still open. {@link #answer} deliberately leaves that stamp in
* place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually
* resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that
* task's future in the same call. So a second task can never reach "{@code turnId} stamped, future
* still open" — the exact pair this method matches on — while a first one already holds it: by
* the time the stamp is gone, so is the eligibility. This is an emergent property of those two
* facts holding together, not something this method (or its callers) enforces on its own — flip
* {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being
* true, with nothing left to fail loudly. The callers below still handle "more than one" as
* defence in depth against exactly that, not because they exercise it today: {@link #reply}
* treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest
* deterministically (its own {@code matching} list has no such guarantee — see its javadoc).
* <p><strong>Can return more than one entry — reachable, not just defence in depth.</strong>
* {@link #send} refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is
* true, and that check matches ANY task whose {@code turnId} is still stamped in
* {@code asyncTasksByTurn}. While a task's {@code turnId} stays stamped — {@link #answer} leaves
* it in place ({@code clearAsyncQuestion(turnId, false)}) until {@link #finishAsyncTask} removes
* the stamp and completes the future in the same call — no second task on the same target can
* reach an eligible state, because {@link #send} would refuse it as BUSY first. That single-task
* guarantee holds only for the {@code turnId}-stamped half of this method's match: an
* {@link Task#askTimedOut} task is, by construction, no longer stamped in {@code asyncTasksByTurn}
* (that is the whole point of forgetting {@code turnId} in {@link #clearAsyncQuestion}), so the
* target is free the moment one ask lapses. A fresh, independent {@code sendAsync} to the same
* target can then be dispatched, itself pause on {@code fleet_ask}, and itself time out — landing
* a second {@code askTimedOut} task on the very target the first one is still waiting to answer
* for. Two (or more) genuinely open tasks on one target is therefore a real, reachable state
* today, not a hypothetical: {@link #reply} treats it as unresolvable and falls back to the
* inbox rather than guess which task a reply belongs to (guessing wrong would hand the lead a
* plausible-looking answer to a delegation the worker never touched — worse than a failure,
* because the lead acts on it); {@link #abandon} instead picks the oldest deterministically (its
* own {@code matching} list has a different, wider match — see its javadoc).
*/
private List<Task> askAnsweredAsyncTasks(String target) {
List<Task> candidates = new ArrayList<>();
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null && task.turnId != null
&& !task.future.isDone()) {
if (target.equals(task.target) && task.question == null && !task.future.isDone()
&& (task.turnId != null || task.askTimedOut)) {
candidates.add(task);
}
}
@@ -898,6 +925,12 @@ public final class MessageService {
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of
// asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is
// what keeps the target from staying BUSY forever), but it would otherwise also erase
// askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still
// belongs to this task, stranding it in the inbox with a false "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
@@ -1162,6 +1195,23 @@ public final class MessageService {
return task;
}
/**
* Mark {@code turnId}'s task as having a {@code fleet_ask} that lapsed with no answer (fleetd #307), so
* {@link #askAnsweredAsyncTasks} still recognizes the worker's eventual real {@code fleet_reply}
* as belonging to it after {@link #clearAsyncQuestion}'s {@code forgetTurn=true} erases
* {@link Task#turnId} — see {@link Task#askTimedOut}. Must be called before that forgetting, while
* {@code turnId} can still resolve the task in {@code asyncTasksByTurn}; a lookup afterward would
* find nothing. Only when it matches the task's current turn — same guard as
* {@link #clearAsyncQuestion} — so a chained second {@code fleet_ask} (#282) that already moved
* the task to a fresh {@code turnId} cannot mark it for a turn that is no longer its own.
*/
private void markAskTimedOut(String turnId) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null && turnId.equals(task.turnId)) {
task.askTimedOut = true;
}
}
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
// CB-582: tell the push loop first — like ticketCollected, a removal for a turnId it never
@@ -1180,14 +1230,61 @@ public final class MessageService {
}
}
/** Complete and detach an async ticket after its worker's actual terminal reply. */
/**
* Complete and detach an async ticket after its worker's actual terminal reply.
*
* <p><strong>fleetd #324.</strong> {@code task.turnId} is read into {@code turnId} exactly once.
* It used to be read twice — once for the null check, once as the removal key — and {@code
* volatile} makes each of those reads individually fresh but does not make the pair atomic.
* {@link #answer} calls this while holding {@code sessionLocks} for the target; {@link #ask}'s
* own timeout path calls {@link #clearAsyncQuestion} (which nulls {@link Task#turnId}) under no
* lock at all. When that unlocked null-out landed between the two reads here, the second read saw
* {@code null} and {@code asyncTasksByTurn.remove(null, task)} threw {@code NullPointerException}
* on the lead's own {@code answer()} call — even though {@code task.future.complete(result)} on
* the line above had already run, so the answer was in fact delivered. Capturing the field once
* removes the torn read; see the ticket for why the wider asymmetry between the locked and
* unlocked sides is not fixed by this alone.
*/
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
String turnId = task.turnId;
if (turnId != null) {
if (finishAsyncTaskRaceHook != null) {
// Test-only (fleetd #324): see the field's own javadoc.
finishAsyncTaskRaceHook.run();
}
asyncTasksByTurn.remove(turnId, task);
}
}
/**
* Null in production; test seam for fleetd #324 — invoked from {@link #finishAsyncTask(Task,
* Reply)} right after {@code task.turnId}'s null-check passes and before the (now-local) value is
* used for the removal. A test installs this to force, deterministically, the exact interleaving
* that a real race between this method and {@link #ask}'s unlocked timeout cleanup can otherwise
* only produce by chance: firing it here reproduces "the field went null between the check and the
* use" against the pre-fix code, and demonstrates the fix tolerates it (the captured local is used
* unconditionally, so a hook that nulls the field afterward cannot affect this call).
*/
private volatile Runnable finishAsyncTaskRaceHook;
/**
* Test-only (fleetd #324): install {@link #finishAsyncTaskRaceHook}. Package-private so the test,
* in the same package, can reach it without widening any production API.
*/
void setFinishAsyncTaskRaceHookForTest(Runnable hook) {
this.finishAsyncTaskRaceHook = hook;
}
/**
* Test-only (fleetd #324): run the exact production cleanup {@link #ask}'s own timeout path runs
* unlocked — {@link #clearAsyncQuestion(String, boolean)} with {@code forgetTurn=true} — so a test
* can reproduce that specific mutation instead of hand-rolling an approximation of it.
*/
void forgetTurnForTest(String turnId) {
clearAsyncQuestion(turnId, true);
}
/** Complete the async ticket correlated to a specific answered turn. */
private void finishAsyncTask(String turnId, Reply result) {
Task task = asyncTasksByTurn.get(turnId);
@@ -5,10 +5,14 @@ import java.util.List;
/**
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps and
* reachability so that a pre-existing config behaves identically after upgrade.
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps
* ({@code maxLoad}) so that a pre-existing config behaves identically after upgrade — capacity
* gating for automatic placement is deliberately out of scope for {@code fixed}, exactly as it
* always has been. Reachability is a narrower exception (fleetd #315, below): a profile is never
* checked for reachability up front, only skipped once it has already failed in <em>this same</em>
* spawn call's retry loop — see the unreachable case below.
*
* <p>Three exceptions walk past the default instead of returning it unconditionally:
* <p>Four exceptions walk past the default instead of returning it unconditionally:
* <ul>
* <li>Quarantine (CB-578 stage B): a quarantined default is a credential that just refused on
* a usage limit, not a transient capacity or reachability concern.
@@ -16,13 +20,21 @@ import java.util.List;
* ({@code BackendOutagePolicy}) — a separate, shorter-lived source from quarantine. When a
* profile is both quarantined and cooling off, only the quarantine reason is reported
* (exhaustion takes priority), matching {@code CompositePeerLauncher}'s explicit-spawn order.
* <li>Unreachable (fleetd #315): {@code CompositePeerLauncher.spawn} retries a failed candidate
* on the next one and rebuilds the {@link PlacementContext} so {@code ctx.unreachable()}
* names every profile that already failed with {@code PeerUnreachableException} in this same
* call. Without this check {@code select} kept handing back the same dead default forever —
* the retry loop's own comment says "so the policy excludes this profile", and this is what
* makes that true for {@code fixed} too, matching {@code weighted}/{@code round-robin}
* (both filter on {@code ctx.unreachable()} via {@link PlacementPolicyUtil#available}).
* <li>Weight 0 (CB-554): {@code fixed} is still automatic selection, so a profile the operator
* marked "never auto-select me" ({@code weight <= 0}) must be skipped here exactly as
* {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming
* the profile is unaffected, only this automatic fallback walk.
* </ul>
* A fleet where nothing is ever quarantined, cooling off, or weight-0 never exercises any of these
* paths, so today's behaviour is unchanged.
* A fleet where nothing is ever quarantined, cooling off, unreachable, or weight-0 never exercises
* any of these paths, so today's behaviour is unchanged — in particular, the very first selection
* of a spawn call always sees an empty {@code unreachable} set, so the first choice is untouched.
*/
final class FixedPlacementPolicy implements PlacementPolicy {
@@ -30,12 +42,12 @@ final class FixedPlacementPolicy implements PlacementPolicy {
public PlacementCandidate select(PlacementContext ctx) {
String d = ctx.defaultProfile();
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d) && !ctx.coolingOff().contains(d)
&& !weightExcluded(ctx, d)) {
&& !ctx.unreachable().contains(d) && !weightExcluded(ctx, d)) {
return new PlacementCandidate(d, null, 1.0f, null);
}
for (PlacementCandidate c : ctx.candidates()) {
if (!ctx.quarantined().contains(c.profile()) && !ctx.coolingOff().contains(c.profile())
&& !c.excluded()) {
&& !ctx.unreachable().contains(c.profile()) && !c.excluded()) {
return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad());
}
}
@@ -44,8 +56,9 @@ final class FixedPlacementPolicy implements PlacementPolicy {
// Exhaustion quarantine takes priority: reported only when quarantine is absent, so the
// message never claims "cooling off" for a profile that is really backend-exhausted.
boolean dCoolingOff = !dQuarantined && ctx.coolingOff().contains(d);
boolean dUnreachable = ctx.unreachable().contains(d);
boolean dWeightExcluded = weightExcluded(ctx, d);
if (dQuarantined || dCoolingOff || dWeightExcluded) {
if (dQuarantined || dCoolingOff || dUnreachable || dWeightExcluded) {
List<String> reasons = new ArrayList<>();
if (dQuarantined) {
reasons.add("is quarantined (backend exhausted)");
@@ -53,6 +66,9 @@ final class FixedPlacementPolicy implements PlacementPolicy {
if (dCoolingOff) {
reasons.add("is cooling off after repeated backend errors");
}
if (dUnreachable) {
reasons.add("is unreachable");
}
if (dWeightExcluded) {
reasons.add("has weight 0 (excluded from automatic selection)");
}
@@ -62,7 +78,7 @@ final class FixedPlacementPolicy implements PlacementPolicy {
}
if (!ctx.candidates().isEmpty()) {
throw new PlacementException("all worker profiles are excluded from automatic "
+ "selection (quarantined, cooling off, or weight-0)");
+ "selection (quarantined, cooling off, unreachable, or weight-0)");
}
throw new PlacementException("no worker profiles configured");
}
@@ -19,6 +19,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.ShuttingDownException;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.WorktreeRequest;
@@ -501,6 +502,11 @@ public final class FleetApp {
ctx.status(201).json(view(member));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (ShuttingDownException e) {
// fleetd #308: the daemon's shutdown drain has already started — 503, not a bare 500,
// so this reads the same as PlacementException below: valid request, refused because
// of a transient daemon state rather than a bad argument.
ctx.status(503).json(Map.of("error", "shutting_down", "detail", e.getMessage()));
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign,
// likely-transient refusal, distinct from "profile does not exist" below. 503: the
@@ -21,6 +21,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import java.util.function.LongSupplier;
@@ -61,6 +62,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
@@ -71,6 +74,16 @@ public final class SessionManager implements TurnListener {
*/
private volatile String fleetRepoRoot;
/**
* fleetd #308: flips true the instant {@link #drainAll} starts, before its registry snapshot
* is even taken — so a spawn already in flight sees the refusal as early as a plain flag can
* make it. This alone cannot close the race completely: a caller that read {@code false} just
* before the flip can still land in the registry after the snapshot. {@link #drainAll}'s
* post-loop sweep is what catches that straggler; the two mechanisms are deliberately paired,
* see {@link #drainAll}'s javadoc.
*/
private final AtomicBoolean draining = new AtomicBoolean(false);
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
/** CB-516: notified with a {@link ReleaseDetail} on every release; no-op until wired. */
@@ -102,13 +115,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;
}
/**
@@ -181,6 +204,14 @@ public final class SessionManager implements TurnListener {
public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest wt,
String sessionName, String resumeSessionId) {
// fleetd #308: refuse before anything else runs — no slot reservation, no launcher spawn —
// so a caller learns the daemon is going down instead of getting a session drainAll will
// never see again. Checked here because every other acquire(...) overload delegates to
// this one, so this is the single point every spawn path passes through.
if (draining.get()) {
throw new ShuttingDownException("fleetd is shutting down; refusing to spawn a session "
+ "the shutdown drain would never see");
}
MemberRole memberRole = (role == null) ? MemberRole.DEV : role;
requireResumeCapability(profile, resumeSessionId);
// CB-619 / fleetd #123: an explicit profile bypasses placement (CompositePeerLauncher only
@@ -272,10 +303,32 @@ 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)) {
// A lifecycle transition replaced the record between the caller's check and this remove.
// Log it: this race is by definition unobservable otherwise, and a reaper that silently
// declines to reap is the hardest kind of behaviour to diagnose after the fact.
log.debug("skipping reap of pane={}: its registry record changed after the idle check "
+ "(most likely a delivery made it BUSY)", expected.paneId());
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) {
@@ -333,6 +386,31 @@ public final class SessionManager implements TurnListener {
// from the registry with no pane stop is an orphaned pane — a live terminal burning a fleet
// slot that no longer appears in the roster and can never be reclaimed.
launcher.stop(paneId);
if (removed != null && !preserveWorktree && removed.worktree() != null) {
// fleetd #316: the `dirty` read above ran while the worker could still write to this
// worktree, so a stale `false` must not be trusted to authorise the --force removal
// below. Re-read the worktree's state one more time, right here — immediately before
// the one step that would destroy it, and only on the path that is actually about to
// do that (invariant 4: no second unconditional `git status` on a release that already
// decided to preserve). By now `launcher.stop` has returned, so this read reflects
// whatever the worker managed to write up to and including its teardown, not whatever
// it had written at release-start time.
if (dirtyImmediatelyBeforeRemoval(removed)) {
preserveWorktree = true;
// The pre-stop snapshot above never ran for this session (the pre-stop read said
// clean), so this is the only chance to get the newly-discovered work into
// refs/wip/* rather than leaving the on-disk preserve as the sole copy. Best-effort,
// like every other snapshot attempt — trySnapshot logs and swallows its own failure.
String lateSnapshotRef = trySnapshot(removed, cause);
log.warn("release {} preserves worktree {} for pane={} terminal={}: it reported "
+ "clean before the pane stopped but dirty immediately before removal — the "
+ "worker wrote to it during teardown, and --force removing it now would "
+ "have destroyed that work{}",
cause, removed.worktree(), paneId, removed.terminalId(),
lateSnapshotRef == null ? "" : " (snapshotted to refs/wip/" + removed.branch()
+ " commit=" + lateSnapshotRef + ")");
}
}
if (removed != null && !preserveWorktree && removed.worktree() != null) {
// fleetd #283: this is the one cleanup step in this method that used to be bare. By the
// time it runs, the registry entry, the retained handle, and the pane are all already
@@ -351,6 +429,24 @@ public final class SessionManager implements TurnListener {
}
}
/**
* fleetd #316: the read that actually authorises {@code worktrees.remove}, taken with the
* worker's pane already stopped. Fails toward preserving (returns {@code true}) on any
* exception — the same rule the pre-stop check applies (CB-581): once we can no longer tell
* whether the worktree is dirty, preserving costs disk while deleting on a guess can destroy
* work that has no other copy.
*/
private boolean dirtyImmediatelyBeforeRemoval(MemberSession removed) {
try {
return worktrees.hasUncommitted(removed.worktree());
} catch (RuntimeException e) {
log.warn("release could not re-check worktree {} for pane={} terminal={} immediately "
+ "before removal; preserving it rather than risk destroying unsaved work: {}",
removed.worktree(), removed.paneId(), removed.terminalId(), e.toString());
return true;
}
}
/**
* Best-effort snapshot of a dirty worktree into {@code refs/wip/<branch>} (CB-578 stage C). A
* failure here must never escalate: the caller has already decided to preserve the worktree
@@ -846,14 +942,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);
@@ -884,10 +984,40 @@ public final class SessionManager implements TurnListener {
* a reason to delete a worker's only copy of its uncommitted work. A session still {@code BUSY}
* when the timeout expired is abandoned mid-turn and logged loudly so an operator can find its
* kept worktree.
*
* <p>fleetd #308: {@code roster()} is a one-shot snapshot (see its javadoc), and nothing used
* to stop a new session from registering after it was taken — {@link #acquire} stayed open for
* as long as this drain waited on a {@code BUSY} session, up to the whole {@code timeoutNanos}
* budget. Two things close that window, deliberately paired because neither alone is complete:
* {@link #draining} is flipped true before the snapshot is even taken, so {@link #acquire}
* refuses (invariant 3: loudly, via {@link ShuttingDownException}) as much of the window as a
* plain flag can close; and the sweep below re-reads the registry once the initial snapshot has
* fully drained and drains whatever a straggler — a caller that read the flag as {@code false}
* a moment before it flipped — still managed to register. The sweep shares the same
* {@code deadline} rather than getting its own: {@code timeoutNanos} is a budget for the WHOLE
* drain (see above), and a straggler must not buy the drain more time than the flag it lost the
* race against would have. In the ordinary case the sweep finds nothing and costs one empty
* {@link #roster()} call.
*/
void drainAll(long timeoutNanos) {
long deadline = System.nanoTime() + timeoutNanos;
for (MemberSession s : roster()) {
draining.set(true);
drainSnapshot(roster(), deadline);
List<MemberSession> stragglers = roster();
if (!stragglers.isEmpty()) {
log.warn("drain sweep found {} session(s) registered after the drain snapshot was "
+ "taken (raced past the shutdown guard); draining them too", stragglers.size());
drainSnapshot(stragglers, deadline);
}
}
/**
* Drain exactly the sessions in {@code snapshot}, waiting out a {@code BUSY} one against the
* shared whole-drain {@code deadline} before releasing it. Shared by {@link #drainAll}'s main
* pass and its post-loop straggler sweep (fleetd #308) so both honor the same one budget.
*/
private void drainSnapshot(List<MemberSession> snapshot, long deadline) {
for (MemberSession s : snapshot) {
try {
if (s.state() == MemberSession.State.BUSY) {
while (System.nanoTime() < deadline) {
@@ -0,0 +1,18 @@
package dev.ltms.fleet.session;
/**
* Thrown by {@link SessionManager#acquire} when a spawn is requested after the daemon's shutdown
* drain has already begun (fleetd #308).
*
* <p>{@link SessionManager#drainAll} snapshots the registry once and tears down exactly what is
* in that snapshot. A session registered after the snapshot is invisible to the drain loop: its
* pane is left running and its worktree is never preserved, and nothing else ever reclaims
* either — the daemon's in-memory registry dies with the process. Refusing the spawn here,
* loudly, is what stops that session from ever being created in the first place, rather than
* silently handing the caller a session the daemon can no longer manage.
*/
public final class ShuttingDownException extends RuntimeException {
public ShuttingDownException(String message) {
super(message);
}
}
@@ -103,6 +103,54 @@ class CallerResolverTest {
assertEquals(Role.PRIMARY, p.role(), "the historical behaviour, now an explicit choice");
}
// ── fleetd #317: an unresolvable caller must never be promoted to the primary ──────────────────
// #305 closed the trigger where a resolved pid matched no pane *and* had no ancestry walk to
// save it. This is the other trigger PaneLocator's javadoc names: the pid never resolves at
// all — LsofPeerPidLookup returns -1 on any failure, including (silently) "lsof found no
// match" — so there is no candidate pid for an ancestry walk to even attempt.
/**
* The failing-without-the-fix case. Before #317's fix, {@code c.terminal() == null} was the
* only test in the loopback-trust fallback, and an unresolved pid produces exactly that same
* {@code null} terminal as a genuine primary — so this caller was handed
* {@code Principal.primary(...)}, a real worker's failed lookup becoming indistinguishable from
* the lead.
*/
@Test
void aFailedPeerPidLookupIsRefusedNotPromotedToPrimary() {
ConnectionIdentity unresolved = new ConnectionIdentity(new PaneLocator(herdr), _ -> -1);
Principal p = new CallerResolver(unresolved).resolve("127.0.0.1", 55555, null);
assertEquals(Role.ANONYMOUS, p.role(),
"an unresolvable caller must never be silently promoted to the primary");
}
/**
* The companion invariant #317 must not break: a caller whose lookup genuinely succeeded, and
* who simply owns no herdr pane — the real primary's own connection — is still the primary.
* This is {@link #loopbackTrustTreatsANonWorkerLoopbackCallerAsThePrimary} pinned again here,
* named for #317 and placed next to the test it must be distinguished from: same {@code null}
* terminal, opposite verdict, because {@code Caller.resolved()} tells them apart.
*/
@Test
void aRealPidThatOwnsNoPaneIsStillThePrimaryNotRefused() {
Principal p = new CallerResolver(nonWorkerIdentity()).resolve("127.0.0.1", 55555, null);
assertEquals(Role.PRIMARY, p.role());
}
/** #317 point 4: token mode never consults {@code c.pid()}, so a failed lookup must not change it. */
@Test
void tokenModeIsUndisturbedByAnUnresolvedLookup() {
ConnectionIdentity unresolved = new ConnectionIdentity(new PaneLocator(herdr), _ -> -1);
CallerResolver r = new CallerResolver(unresolved, true, "s3cret");
assertEquals(Role.ANONYMOUS, r.resolve("127.0.0.1", 55555, null).role(),
"no credential is still just ANONYMOUS, as before #317 — unchanged by the lookup failing");
assertEquals(Role.PRIMARY, r.resolve("127.0.0.1", 55555, "Bearer s3cret").role(),
"a valid token still authenticates the primary even though the peer-pid lookup failed");
}
@Test
void tokenModeRefusesANonWorkerCallerThatPresentsNoToken() {
Principal p = new CallerResolver(nonWorkerIdentity(), true, "s3cret")
@@ -43,6 +43,22 @@ class ConnectionIdentityTest {
assertNull(with(_ -> 999_999).callerTerminal("127.0.0.1", 55555));
}
@Test
void callerIsUnresolvedWhenThePeerPidLookupFails() {
// fleetd #317: LsofPeerPidLookup returns -1 on any failure — a fork error, or (silently)
// simply no matching lsof line. Caller.resolved() is the one place that sentinel is tested.
ConnectionIdentity.Caller c = with(_ -> -1).resolve("127.0.0.1", 55555);
assertFalse(c.resolved(), "a -1 pid means the lookup failed, not that this pid owns no pane");
}
@Test
void callerIsResolvedWhenThePidIsRealEvenThoughItOwnsNoPane() {
// The primary's own connection: a real, lsof-found pid that just isn't a worker pane. This
// must read as "resolved" — the distinction #317 turns on.
ConnectionIdentity.Caller c = with(_ -> 999_999).resolve("127.0.0.1", 55555);
assertTrue(c.resolved());
}
@Test
void resolvesTheCallersPidAndCwd() {
// CB-112: the primary maps to no pane, but its PID and cwd are still readable.
@@ -189,7 +189,7 @@ class FleetMcpTest {
}
@Test
void unansweredAsyncAskReturnsTheTicketToPending() throws Exception {
void unansweredAsyncAskReturnsTheTicketToPendingThenAWorkersLateReplyStillCompletesIt() throws Exception {
McpSchema.CallToolResult accepted = FleetMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
@@ -203,8 +203,15 @@ class FleetMcpTest {
assertTrue(textOf(ask).contains("no answer"), textOf(ask));
assertTrue(textOf(FleetMcp.poll(messages, ticket, null)).startsWith("[pending"));
// fleetd #307: the worker resumed on its own after the primary never answered, and its real
// fleet_reply must complete its OWN async ticket — not strand in the inbox with
// fleet_poll{ticket} stuck PENDING forever and later force-failed with a false "session
// released before it replied" reason. This used to land in the inbox instead (see the old
// assertion this replaced: messages.drainReplies("term_a").getFirst()...) — that was the bug.
FleetMcp.reply(messages, "term_a", "finished after timeout");
assertEquals("finished after timeout", messages.drainReplies("term_a").getFirst().content());
assertEquals("finished after timeout", textOf(FleetMcp.poll(messages, ticket, null)));
assertTrue(messages.drainReplies("term_a").isEmpty(),
"the reply completed its own ticket directly and never touched the inbox");
}
@Test
@@ -615,6 +615,67 @@ class CompositePeerLauncherTest {
assertEquals(1, adapter.spawnCount("b"));
}
/**
* fleetd #315: {@code CompositePeerLauncher.spawn} rebuilds the {@link PlacementContext} after
* every failed attempt "so the policy excludes this profile" (see the comment at the retry call
* site) — but {@code FixedPlacementPolicy} never read {@code ctx.unreachable()}, so under the
* default {@code fixed} placement every retry re-picked the same dead default and a second,
* healthy, configured profile was never tried. This is the same scenario as
* {@link #failoverRetriesNextCandidateWhenProfileIsUnreachable}, but pinned to {@code fixed()}
* instead of {@code weighted()} — the three existing failover tests all use {@code weighted()},
* which is exactly why nobody caught this: the retry loop's contract has no coverage under its
* own default policy.
*/
@Test
void failoverRetriesNextCandidateUnderFixedPlacementWhenProfileIsUnreachable() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"a", stubWorker("a"),
"b", stubWorker("b"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of("a"));
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "a", profiles, PlacementPolicies.fixed(), _ -> 0);
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
assertEquals("b", h.profile(),
"fixed placement must fail over from the unreachable default a to the healthy b");
assertEquals(1, adapter.spawnCount("a"), "a was tried once and failed");
assertEquals(1, adapter.spawnCount("b"), "b was tried once and succeeded");
}
/**
* fleetd #315: the same fix — {@code FixedPlacementPolicy} consulting {@code ctx.unreachable()}
* — also covers the wiring-bug branch in {@code CompositePeerLauncher.spawn}: a profile that
* placement is allowed to choose (it is in the configured candidate list) but that no delegate
* declares ({@code byProfile.get(chosen.profile()) == null}). That branch adds the profile to
* {@code unreachable} and {@code continue}s without ever calling a launcher, so before this fix
* {@code fixed} handed back the same adapterless profile on every remaining attempt too.
*/
@Test
void failoverSkipsAConfiguredProfileNoAdapterDeclaresUnderFixedPlacement() {
FakeHerdr herdr = new FakeHerdr();
// Placement's candidate list has three profiles, in this order (LinkedHashMap preserves it,
// and the fixed default resolves to the first — see the `ordered` helper's own javadoc).
Map<String, FleetConfig.Profile> profiles = new LinkedHashMap<>();
profiles.put("c", stubWorker("c"));
profiles.put("a", stubWorker("a"));
profiles.put("b", stubWorker("b"));
// The adapter only declares a and b — c is a configured profile with no owning adapter,
// the "wiring bug" the comment in CompositePeerLauncher.spawn calls out.
Map<String, FleetConfig.Profile> adapterProfiles = new LinkedHashMap<>();
adapterProfiles.put("a", stubWorker("a"));
adapterProfiles.put("b", stubWorker("b"));
StubLauncher adapter = new StubLauncher("claude", herdr, adapterProfiles, "a", Set.of());
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "a", profiles, PlacementPolicies.fixed(), _ -> 0);
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
assertEquals("a", h.profile(),
"c has no adapter, so fixed placement must skip it and land on the next candidate, a");
assertEquals(0, adapter.spawnCount("c"), "c is never spawned — no adapter owns it");
assertEquals(1, adapter.spawnCount("a"));
}
@Test
void explicitSpawnAtMaxLoadThrowsPlacementExceptionNamingProfileLiveAndCap() {
FakeHerdr herdr = new FakeHerdr();
@@ -0,0 +1,264 @@
package dev.ltms.fleet.msg;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Delivery;
import com.rabbitmq.client.Envelope;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Proxy;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-318: a delivery landing on the consumer work-pool thread <em>after</em>
* {@link AmqpReplyInbox#release} has already swapped the target's {@code held} entry for its
* tombstone, but <em>before</em> {@code release()} itself returns, must be nacked-with-requeue —
* never silently retained in a fresh map {@code release()} has already stopped looking at.
*
* <p><strong>This forces the actual interleaving, not a sequence of calls.</strong> {@code release()}
* runs on its own thread and is made to block <em>inside</em> its nack loop, via a fake
* {@link Channel} whose {@code basicNack} blocks on its first invocation. That block is only
* reachable after {@code release()}'s {@code held.compute(...)} has already swapped in the
* {@code RELEASED} tombstone (the compute call happens strictly before the loop that calls
* {@code basicNack}), so observing it is direct, ordering-guaranteed proof that the tombstone is in
* place and {@code release()} has not yet returned — still holding {@code channelLock} — when a
* second thread fires {@code own()}'s captured {@link DeliverCallback} for a brand-new message on the
* same target. No mocking library is on the classpath, so the fake broker is a {@link Proxy}, the
* same pattern {@code AmqpReplyInboxRecoveryRaceTest} already uses.
*
* <p><strong>What this does and does not prove.</strong> It proves that a delivery whose
* {@code computeIfAbsent} call is ordered strictly after {@code release()}'s tombstone swap — while
* {@code release()} is still running — is nacked-with-requeue rather than silently parked forever.
* It does not drive a real broker: {@code basicNack} here is a recorded call on a fake channel, not a
* verified requeue-and-redeliver. That half of the contract (a nacked-with-requeue delivery really
* does come back to a later owner) is already covered against a real broker by
* {@code AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery}, which
* this test does not duplicate.
*/
class AmqpReplyInboxReleaseRaceTest {
@Test
@Timeout(15)
void deliveryArrivingWhileReleaseIsStillRunningIsNackedNotStranded() throws Exception {
String target = "worker-release-race";
List<long[]> nacks = new CopyOnWriteArrayList<>(); // {deliveryTag, requeue(1/0)}
List<Long> acks = new CopyOnWriteArrayList<>();
AtomicReference<DeliverCallback> deliverCallback = new AtomicReference<>();
CountDownLatch nackStarted = new CountDownLatch(1);
CountDownLatch releaseMayFinishNack = new CountDownLatch(1);
AtomicInteger nackCallCount = new AtomicInteger();
Channel consumeChannel = fakeConsumeChannel(deliverCallback, nacks, acks, nackCallCount,
nackStarted, releaseMayFinishNack);
Channel publishChannel = fakeInertChannel();
Connection connection = fakeConnection(consumeChannel, publishChannel);
AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH);
inbox.own(target);
assertTrue(deliverCallback.get() != null, "own() must have registered a DeliverCallback");
// Seed one already-held delivery (m0) so release()'s nack loop has something to iterate, and
// therefore somewhere to block, before it can return.
deliverCallback.get().handle("ctag", delivery(1L, "m0", "first"));
AtomicReference<Throwable> releaseError = new AtomicReference<>();
Thread releaseThread = new Thread(() -> {
try {
inbox.release(target);
} catch (Throwable t) {
releaseError.set(t);
}
}, "release-under-test");
releaseThread.start();
// This latch only fires from inside the fake channel's basicNack — i.e. from inside
// release()'s nack loop, which release()'s code only reaches AFTER held.compute(...) has
// already swapped in RELEASED. Waiting for it is direct proof the swap has happened and
// release() has not yet returned (it is stuck mid-loop, still holding channelLock).
assertTrue(nackStarted.await(10, TimeUnit.SECONDS),
"release() never reached its nack loop — it may not have started");
// The exact interleaving CB-318 describes: a delivery for a NEW message on the same target
// lands on the "consumer work-pool thread" (this second thread) while release() is still
// running. With the pre-fix code (a bare held.remove(target)) this created a brand-new map
// under computeIfAbsent that release() — already past its remove — never looks at again.
AtomicReference<Throwable> deliveryError = new AtomicReference<>();
Thread deliveryThread = new Thread(() -> {
try {
deliverCallback.get().handle("ctag", delivery(2L, "m1", "second"));
} catch (Throwable t) {
deliveryError.set(t);
}
}, "concurrent-delivery");
deliveryThread.start();
// Head start for the delivery thread to reach (and, on the fixed code, block on)
// channelLock — release() still holds it at this point, so a correct fix cannot have
// resolved m1's nack yet. Purely in-memory work (computeIfAbsent, a reference compare)
// separates deliveryThread.start() from that block point, so 300ms is a large margin, not a
// tight timing assumption.
Thread.sleep(300);
assertEquals(1, nackCallCount.get(),
"the concurrent delivery must not resolve its nack before release() gives up "
+ "channelLock — if this is 2 already, the interleaving below is not being "
+ "tested, only a sequential call");
releaseMayFinishNack.countDown(); // let release() finish nacking m0 and return
assertTrue(releaseThread.join(Duration.ofSeconds(10)), "release() did not finish");
assertTrue(deliveryThread.join(Duration.ofSeconds(10)), "the concurrent delivery did not finish");
assertNull(releaseError.get(), "release() threw: " + releaseError.get());
assertNull(deliveryError.get(), "the concurrent delivery threw: " + deliveryError.get());
assertEquals(2, nacks.size(),
"both the pre-held m0 and the concurrently-arriving m1 must be nacked, got: "
+ nacks.stream().map(n -> "[tag=" + n[0] + " requeue=" + n[1] + "]").toList());
assertTrue(nacks.stream().allMatch(n -> n[1] == 1L),
"invariant 1 (never drop): every nack must set requeue=true");
assertTrue(nacks.stream().anyMatch(n -> n[0] == 1L), "m0's delivery tag must be nacked");
assertTrue(nacks.stream().anyMatch(n -> n[0] == 2L),
"m1 — delivered while release() was still running, after the tombstone swap — must be "
+ "nacked, not silently retained in a map release() will never look at again");
assertTrue(acks.isEmpty(), "invariant 1 (never drop): a held reply must never be basicAck'd");
}
private static Delivery delivery(long tag, String msgId, String body) {
Envelope envelope = new Envelope(tag, false, "", "irrelevant");
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder().messageId(msgId).build();
return new Delivery(envelope, props, body.getBytes(StandardCharsets.UTF_8));
}
/** A {@link Proxy}-backed consume {@link Channel}: blocks the FIRST {@code basicNack} call on
* {@code releaseMayFinishNack}, after signalling {@code nackStarted} — everything else records
* the call and returns a harmless default, matching the style already used by
* {@code AmqpReplyInboxRecoveryRaceTest}. */
private static Channel fakeConsumeChannel(AtomicReference<DeliverCallback> deliverCallback,
List<long[]> nacks, List<Long> acks,
AtomicInteger nackCallCount,
CountDownLatch nackStarted,
CountDownLatch releaseMayFinishNack) {
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("basicConsume")) {
deliverCallback.set((DeliverCallback) args[2]);
return "ctag";
}
if (name.equals("basicNack")) {
long tag = (long) args[0];
boolean requeue = (boolean) args[2];
if (nackCallCount.incrementAndGet() == 1) {
nackStarted.countDown();
if (!releaseMayFinishNack.await(10, TimeUnit.SECONDS)) {
throw new IllegalStateException("test never released the nack latch");
}
}
nacks.add(new long[] {tag, requeue ? 1L : 0L});
return null;
}
if (name.equals("basicAck")) {
acks.add((long) args[0]);
return null;
}
if (name.equals("equals")) {
return proxy == args[0];
}
if (name.equals("hashCode")) {
return System.identityHashCode(proxy);
}
if (name.equals("toString")) {
return "FakeConsumeChannel";
}
return defaultValue(method.getReturnType());
};
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
new Class<?>[] {Channel.class}, handler);
}
/** A {@link Proxy}-backed {@link Channel} that answers every call with a harmless default — used
* as the publish channel, which this test never actually publishes on. */
private static Channel fakeInertChannel() {
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("equals")) {
return proxy == args[0];
}
if (name.equals("hashCode")) {
return System.identityHashCode(proxy);
}
if (name.equals("toString")) {
return "FakeInertChannel";
}
return defaultValue(method.getReturnType());
};
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
new Class<?>[] {Channel.class}, handler);
}
/** A {@link Proxy}-backed {@link Connection} handing out {@code first} then {@code second} from
* successive {@code createChannel()} calls, matching {@link AmqpReplyInbox}'s constructor. */
private static Connection fakeConnection(Channel first, Channel second) {
AtomicInteger calls = new AtomicInteger();
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("createChannel") && (args == null || args.length == 0)) {
return calls.getAndIncrement() == 0 ? first : second;
}
if (name.equals("equals")) {
return proxy == args[0];
}
if (name.equals("hashCode")) {
return System.identityHashCode(proxy);
}
if (name.equals("toString")) {
return "FakeConnection";
}
return defaultValue(method.getReturnType());
};
return (Connection) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
new Class<?>[] {Connection.class}, handler);
}
private static Object defaultValue(Class<?> type) {
if (!type.isPrimitive() || type == void.class) {
return null;
}
if (type == boolean.class) {
return Boolean.FALSE;
}
if (type == long.class) {
return 0L;
}
if (type == short.class) {
return (short) 0;
}
if (type == byte.class) {
return (byte) 0;
}
if (type == char.class) {
return (char) 0;
}
if (type == double.class) {
return 0.0d;
}
if (type == float.class) {
return 0.0f;
}
return 0;
}
}
@@ -940,6 +940,54 @@ class MessageServiceTest {
assertEquals("PR opened: https://example/pulls/42", view.reply());
}
/**
* fleetd #324: {@code answer()} holds {@code sessionLocks} for the target and, once the worker's
* real terminal reply arrives, calls {@code finishAsyncTask}, which used to read the volatile
* {@code task.turnId} twice — once to check it is non-null, once as the key for
* {@code asyncTasksByTurn.remove}. {@code ask()}'s own timeout path mutates the same field with no
* lock at all. This test does not wait for a real race to land in that narrow window between the
* two reads — instead it drives the exact sequence the ticket describes (worker asks, primary
* answers, worker's real reply arrives) and, via a package-private test hook wired to fire at
* precisely that point, runs the identical production cleanup {@code ask()}'s timeout catch block
* runs ({@code clearAsyncQuestion(turnId, true)}) so the field goes {@code null} between the two
* reads deterministically rather than by chance.
*
* <p>What this proves: given that exact interleaving, {@code answer()} must not throw and the
* ticket must still resolve to the worker's real reply. What it does not prove: that the
* interleaving itself is reachable in production — that is established by reading the code (see
* the ticket), not by this test, since forcing it via a hook is not the same as two independent
* threads racing on their own schedules.
*/
@Test
void finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
// Fire ask()'s own unlocked timeout cleanup at the moment finishAsyncTask has already checked
// task.turnId is non-null but has not yet used it — the exact torn-read window fleetd #324
// describes.
messages.setFinishAsyncTaskRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the lead's own answer() call must not throw because ask()'s timeout cleanup raced it");
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"the ticket must still resolve to the worker's real reply despite the forced race");
}
@Test
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
@@ -957,6 +1005,76 @@ class MessageServiceTest {
assertEquals(MessageService.Phase.DONE, awaitTicketPhase(next, MessageService.Phase.DONE).phase());
}
/**
* fleetd #307: a worker's {@code fleet_ask} can time out because the primary never answers —
* distinct from {@link #aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket}, where the
* primary DID answer and only its own bounded wait for the resumed turn expired.
* {@code ask()}'s timeout path deliberately forgets the task's {@code turnId} (so
* {@code hasAsyncQuestion} stops reporting the target BUSY — see
* {@code unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget} above), which used
* to also erase the one signal {@code askAnsweredAsyncTasks} needed to recognize the worker's
* eventual real {@code fleet_reply}. That reply then had nowhere to land but the inbox, and
* {@code fleet_poll{ticket}} stayed PENDING forever — later force-failed with the false reason
* "session released before it replied", even though the worker had, in fact, replied.
*/
@Test
void aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks then finishes alone");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT,
messages.ask(T, "which config?", 200).outcome());
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
"only the question wait ended; the delegated turn may still finish");
// The worker keeps working past the timeout and only now calls fleet_reply — with no live
// rendezvous waiter open (ask()'s timeout already closed it) and no new send() having
// reopened one for this target.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"fleet_poll{ticket} must return the worker's real reply, not stay pending forever");
assertEquals("reply", done.replySource());
assertFalse(messages.hasStrandedReply(T),
"the reply completed its own ticket directly and never touched the inbox");
}
/**
* fleetd #307's ambiguity guard: an ask timeout frees its target ({@code hasAsyncQuestion}
* becomes false the instant it lapses — proven above), so a second, independent delegation can
* be dispatched to the same target and itself go on to ask-and-lapse before the first worker's
* real reply ever arrives. Two open tasks are then both eligible candidates on one target with
* no live waiter to disambiguate them. A reply arriving now must not guess which one it answers
* — guessing wrong would hand the lead a plausible-looking answer to a delegation the worker
* never touched, worse than a failure because the lead acts on it — so it must fall back to the
* inbox exactly as the zero-candidate case does.
*/
@Test
void twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess() throws Exception {
String ticket1 = messages.sendAsync(T, "first task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q1?", 200).outcome());
String ticket2 = messages.sendAsync(T, "second task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome());
assertTrue(messages.reply(T, "which task does this answer?"));
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(),
"an ambiguous reply must not guess ticket1");
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket2).phase(),
"an ambiguous reply must not guess ticket2");
assertTrue(messages.hasStrandedReply(T));
var drained = messages.drainReplies(T);
assertEquals(1, drained.size());
assertEquals("which task does this answer?", drained.get(0).content());
}
@Test
void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception {
String first = messages.sendAsync(T, "first task");
@@ -26,7 +26,12 @@ import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.*;
@@ -81,17 +86,50 @@ class SessionManagerTest {
private volatile RuntimeException hasUncommittedFailure;
private volatile RuntimeException snapshotFailure;
private final java.util.concurrent.atomic.AtomicLong snapshotSeq = new java.util.concurrent.atomic.AtomicLong();
/** fleetd #316: successive {@code hasUncommitted} answers, one per call, last one sticky
* once exhausted — models a worktree whose state changes between reads. Empty (the
* default) falls back to the plain {@link #dirty} flag, so every existing test using this
* fake keeps returning one fixed answer. */
private final List<Boolean> dirtySequence = new java.util.concurrent.CopyOnWriteArrayList<>();
private int failHasUncommittedOnCall = -1;
private RuntimeException hasUncommittedCallFailure;
private final java.util.concurrent.atomic.AtomicInteger hasUncommittedCalls =
new java.util.concurrent.atomic.AtomicInteger();
RecordingWorktrees dirty(boolean dirty) {
this.dirty = dirty;
return this;
}
/** fleetd #316: return {@code answers[0]} on the first {@code hasUncommitted} call,
* {@code answers[1]} on the second, and so on; the last element repeats after that. */
RecordingWorktrees dirtySequence(boolean... answers) {
for (boolean a : answers) {
dirtySequence.add(a);
}
return this;
}
int hasUncommittedCallCount() {
return hasUncommittedCalls.get();
}
RecordingWorktrees failHasUncommittedWith(RuntimeException e) {
this.hasUncommittedFailure = e;
return this;
}
/**
* Throw from {@code hasUncommitted} on one specific call only, counting from 0. The
* whole-double {@link #failHasUncommittedWith} cannot express fleetd #316's fail-safe
* case, which needs the pre-stop read to succeed and only the late read to fail.
*/
RecordingWorktrees failHasUncommittedOnCall(int call, RuntimeException e) {
this.failHasUncommittedOnCall = call;
this.hasUncommittedCallFailure = e;
return this;
}
RecordingWorktrees failRemoveFor(String worktreePath) {
failRemoveFor.add(worktreePath);
return this;
@@ -122,9 +160,16 @@ class SessionManagerTest {
@Override
public boolean hasUncommitted(String worktreePath) {
int call = hasUncommittedCalls.getAndIncrement();
if (hasUncommittedFailure != null) {
throw hasUncommittedFailure;
}
if (call == failHasUncommittedOnCall) {
throw hasUncommittedCallFailure;
}
if (!dirtySequence.isEmpty()) {
return dirtySequence.get(Math.min(call, dirtySequence.size() - 1));
}
return dirty;
}
@@ -178,14 +223,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 +720,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};
@@ -824,6 +892,186 @@ class SessionManagerTest {
.count();
}
// --- fleetd #308: a spawn accepted while the shutdown drain is running must not orphan ---
@Test
void acquireRefusesANewSpawnOnceDrainAllHasStarted() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(50)); // empty roster — returns immediately,
// but the shutdown guard it flips must stay tripped for the life of the process.
ShuttingDownException e = assertThrows(ShuttingDownException.class,
() -> sessions.acquire("ltms-local", null, "/caller", "term_primary"),
"a spawn requested after the drain has begun must be refused loudly (invariant 3), "
+ "not silently registered into a registry the drain will never revisit");
assertNotNull(e.getMessage());
assertFalse(e.getMessage().isBlank(), "the refusal must say why, not just that it failed");
assertTrue(sessions.roster().isEmpty(), "the refused spawn must never reach the registry");
}
/**
* fleetd #308: the guard above closes most of the shutdown-race window, but it cannot close
* all of it — a caller that already passed the {@code draining} check before {@code drainAll}
* flips it can still be mid-{@code launcher.spawn()} (a real herdr round trip, not
* instantaneous) when {@code drainAll} takes its registry snapshot. This test forces exactly
* that interleaving with a launcher double that blocks the second {@code spawn()} call and the
* first {@code stop()} call until released, then proves the post-loop sweep in {@code
* drainAll} still finds and tears down the straggler that lands in the registry afterward.
*/
@Test
void drainAllSweepsAStragglerThatRegisteredAfterTheInitialSnapshot() throws Exception {
FakeHerdr herdr = new FakeHerdr();
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 delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
RaceLauncher race = new RaceLauncher(delegate);
SessionManager sessions = new SessionManager(race);
// Registered normally, before the drain starts — the first spawn call, never blocked.
MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
ExecutorService exec = Executors.newFixedThreadPool(2);
try {
// The straggler's acquire() reads `draining == false` (checked before this call ever
// touches the launcher) and then blocks inside its own spawn() — the second spawn call.
Future<MemberSession> straggler = exec.submit(() ->
sessions.acquire("ltms-local", "/late", "/caller", "ownerLate"));
assertTrue(race.enteredSecondSpawn.await(5, TimeUnit.SECONDS),
"the straggler must have passed the shutdown guard and reached spawn() before "
+ "drainAll ever runs");
assertEquals(1, sessions.roster().size(),
"the straggler is still inside spawn() — not registered yet");
// drainAll flips `draining`, snapshots the registry (only `ready` is in it), and starts
// releasing that snapshot — its first release() call stops `ready`'s pane, which this
// launcher double blocks on so the interleaving below is deterministic, not a timing bet.
Future<?> drain = exec.submit(() -> sessions.drainAll(TimeUnit.SECONDS.toNanos(5)));
assertTrue(race.enteredFirstStop.await(5, TimeUnit.SECONDS),
"drainAll must be stopping the ready session's pane — proof its initial "
+ "registry snapshot has already been taken");
// Only now does the straggler's spawn complete and register — strictly after the
// snapshot drainAll's main pass is working from.
race.releaseSecondSpawn.countDown();
MemberSession registered = straggler.get(5, TimeUnit.SECONDS);
// Let drainAll finish releasing `ready`; it then re-checks the registry and must find
// (and drain) the straggler that just landed in it.
race.releaseFirstStop.countDown();
drain.get(5, TimeUnit.SECONDS);
assertTrue(sessions.roster().isEmpty(),
"the post-loop sweep must drain the straggler too, not just the initial snapshot");
assertNotNull(registered.paneId());
long paneCloseCalls = herdr.calls.stream().filter(c -> "pane.close".equals(c.method())).count();
assertEquals(2, paneCloseCalls,
"both ready's pane AND the straggler's pane must actually be stopped — a pane "
+ "left running is exactly the orphan this ticket is about");
} finally {
exec.shutdownNow();
}
}
/**
* Delegates every call while blocking the SECOND {@code spawn()} call and the FIRST
* {@code stop()} call until the test releases them — used to force the fleetd #308 race
* deterministically instead of betting on real thread-scheduling timing.
*/
private static final class RaceLauncher implements PeerLauncher {
private final PeerLauncher delegate;
private final AtomicInteger spawnCalls = new AtomicInteger();
private final AtomicInteger stopCalls = new AtomicInteger();
final CountDownLatch enteredSecondSpawn = new CountDownLatch(1);
final CountDownLatch releaseSecondSpawn = new CountDownLatch(1);
final CountDownLatch enteredFirstStop = new CountDownLatch(1);
final CountDownLatch releaseFirstStop = new CountDownLatch(1);
RaceLauncher(PeerLauncher delegate) {
this.delegate = delegate;
}
private static void awaitOrFail(CountDownLatch latch) {
try {
if (!latch.await(5, TimeUnit.SECONDS)) {
throw new AssertionError("RaceLauncher latch timed out");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new AssertionError("RaceLauncher latch interrupted", e);
}
}
@Override
public Set<Capability> capabilities() {
return delegate.capabilities();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return delegate.capabilitiesFor(profileName);
}
@Override
public PeerHandle spawn(SpawnRequest req) {
if (spawnCalls.incrementAndGet() == 2) {
enteredSecondSpawn.countDown();
awaitOrFail(releaseSecondSpawn);
}
return delegate.spawn(req);
}
@Override
public Set<String> profiles() {
return delegate.profiles();
}
@Override
public String defaultProfile() {
return delegate.defaultProfile();
}
@Override
public String effectiveCwd(SpawnRequest req) {
return delegate.effectiveCwd(req);
}
@Override
public List<String> parityOverlay(String profileName) {
return delegate.parityOverlay(profileName);
}
@Override
public List<?> list() {
return delegate.list();
}
@Override
public int reapOrphanWorkers() {
return delegate.reapOrphanWorkers();
}
@Override
public void stop(String id) {
if (stopCalls.incrementAndGet() == 1) {
enteredFirstStop.countDown();
awaitOrFail(releaseFirstStop);
}
delegate.stop(id);
}
@Override
public boolean clearContext(String id) {
return delegate.clearContext(id);
}
}
private static List<String> promptTexts(FakeHerdr herdr) {
return herdr.calls.stream()
.filter(c -> "agent.prompt".equals(c.method()))
@@ -999,6 +1247,104 @@ class SessionManagerTest {
+ "dirty check threw");
}
// --- fleetd #316: the dirty check must be re-taken after the worker is stopped, not trusted
// stale from before it ------------------------------------------------------------------------
@Test
void releaseDoesNotRemoveAWorktreeThatBecameDirtyBetweenTheFirstCheckAndRemoval() {
// Models the exact race #316 reports: hasUncommitted answers clean while the worker is
// still running (call 1), the worker then writes new work, and by the time release is
// about to force-remove the worktree a second read (call 2) would see it as dirty. Without
// the fix this test fails: release() never re-reads and force-removes the worktree anyway.
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees().dirtySequence(false, true);
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-316a", null));
sessions.release(s.paneId());
assertTrue(worktrees.removeCalls().isEmpty(),
"a worktree that turned dirty between the pre-stop read and removal must be preserved");
assertEquals(2, worktrees.hasUncommittedCallCount(),
"the fix re-reads hasUncommitted exactly once more, immediately before removal");
}
@Test
void releasePreservesAWorktreeWhoseLateRecheckCannotBeRead() {
// fleetd #316 invariant 1, which no test pinned when the fix landed: the late re-check
// fails toward PRESERVING. Found by mutation — flipping dirtyImmediatelyBeforeRemoval's
// catch from `return true` to `return false` turned the guard into a cause of the very
// data loss it was added to stop, and the whole suite stayed green. The pre-stop read
// succeeds and says clean (call 0); the read that authorises the removal throws (call 1).
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees()
.dirtySequence(false)
.failHasUncommittedOnCall(1, new WorktreeException("git status exited 128"));
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-316c", null));
sessions.release(s.paneId());
assertTrue(worktrees.removeCalls().isEmpty(),
"a worktree whose state cannot be read immediately before removal must be kept: "
+ "preserving costs disk, deleting on a guess destroys work with no other copy");
assertEquals(2, worktrees.hasUncommittedCallCount(),
"the late re-check still runs — it is the throwing call, not a skipped one");
}
@Test
void releaseSnapshotsWorkFoundOnlyByTheLateRecheck() {
// #316's second half: the pre-stop dirty=false means trySnapshot never ran for this
// session, so the late-discovered work would otherwise have no refs/wip/* copy at all —
// only the on-disk preserve. The re-check path must snapshot it too.
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees().dirtySequence(false, true);
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-316b", null));
sessions.release(s.paneId());
assertEquals(java.util.List.of(s.worktree()), worktrees.snapshotCalls(),
"the newly-dirty worktree is snapshotted even though the pre-stop check saw it clean");
}
@Test
void releaseStillRemovesAWorktreeThatStaysCleanOnTheLateRecheck() {
// The ordinary, non-racing case: nothing else changes behaviour when the second read
// agrees with the first.
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(false);
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-316c", null));
sessions.release(s.paneId());
assertEquals(java.util.List.of(s.worktree()), worktrees.removeCalls(),
"a worktree that is still clean on the late recheck is removed as before");
}
@Test
void releaseNeverReChecksAWorktreeAlreadyPreservedByTheFirstDirtyCheck() {
// Invariant 4 from #316: no second unconditional git status. A release that already
// decided to preserve (the ordinary CB-576 dirty path) must not pay for a second read.
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(true);
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-316d", null));
sessions.release(s.paneId());
assertEquals(1, worktrees.hasUncommittedCallCount(),
"a release that already preserves on the first read must not re-check before "
+ "skipping the removal it was never going to do");
assertTrue(worktrees.removeCalls().isEmpty());
}
/**
* fleetd #283 defect 1 changed this test's own premise, so its assertions are updated along
* with the production fix. Before #283, the middle session's worktree-removal failure escaped