Compare commits

..

27 Commits

Author SHA1 Message Date
Dai Ha 9d653e86df fleetd #726: name RELAUNCH_NEVER_READY in the IN_PROGRESS terminal-state list
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 2m7s
The javadoc on RollState.IN_PROGRESS enumerates the terminal states the entry
can be overwritten with, and omitted RELAUNCH_NEVER_READY. That state is
reachable at LeadRollover.java:710, so the list told a reader a state could not
occur when it can. Found by a reviewer on PR #742, outside its assigned scope.

Comment only; no behaviour change.
2026-10-04 21:44:03 +02:00
Dai Ha f6d1131d7a Merge remote-tracking branch 'origin/worker/726-unit2-75cb13-4' 2026-10-04 21:43:34 +02:00
Dai Ha a0505dc614 fleetd #726 unit 2: scope the member-daemon assertion to the roll itself
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 2m0s
The assembled daemon's own boot-time orphan-worker reap makes a real call on
the member fake before any roll starts. Clear the member fake's recorded
calls once assembly finishes and before the roll begins, so the assertion
measures calls made since the roll started rather than the whole process's
lifetime, and reword its message to say so.
2026-10-04 21:24:49 +02:00
Dai Ha d2f30f1654 fleetd #726 unit 2: cover the bootstrapText relaunch send with a dedicated test
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 1m46s
LeadRollover already sends every call through the lead-bound AgentControl and
WorkspaceControl it receives at construction, so no production code needed a
routing fix. Add a regression test that drives a full relaunch to the point
where recognition times out and asserts bootstrapText still lands on the lead
daemon and never on the member daemon, the one path the existing assembly test
never reaches.
2026-10-04 21:17:15 +02:00
Dai Ha cd1f04cbb4 Merge remote-tracking branch 'origin/worker/737-owner-key-ff061f-10'
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 53s
CI / build (push) Failing after 1m48s
2026-10-04 21:11:20 +02:00
Dai Ha 73137f198f Merge remote-tracking branch 'origin/main' into worker/726-unit2-75cb13-4 2026-10-04 20:55:29 +02:00
Dai Ha 4ffe49f3bb fleetd #726 unit 2: replace /clear-based lead rollover with a real process restart
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 59s
CI / build (pull_request) Failing after 1m56s
LeadRollover's deferred continuation now ends the old lead's pane, relaunches
a fresh one, and bootstraps it, instead of sending /clear into the same
process. The relaunch step runs two separate bounded waits instead of one
combined check: a readiness wait (the fresh pane reaches a real turn
boundary) is the safety gate and withholds bootstrapText on timeout
(RELAUNCH_NEVER_READY); a recognition wait (the fresh terminal shows up in
the live-lead map) is bookkeeping only, so a timeout there still lets
bootstrapText go out (RELAUNCH_NOT_RECOGNISED). clearSettleSeconds is
retired in favor of relaunchReadySeconds (default 45), which bounds both
waits. Updates FleetConfig/FleetMcp operator-facing text to match.
2026-10-04 20:44:42 +02:00
Dai Ha efd9cdb983 fleetd #737 units 1+2: key tickets and turns on a stable owner, not a terminal
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Failing after 2m4s
Add Principal.ownerKey(): role-prefixed, keyed on name for a named lead and
a collaborator (survives a handover's terminal change), on terminal for a
worker, architect and observer, and on a distinct "anonymous" value for an
unauthenticated caller so that case no longer relies on Authz refusing it
first. The unnamed primary keeps a null key, preserving its primary-wide
ticket rule.

Thread that key through Task.creatorOwner, Rendezvous.Owner, poll,
pendingAsk and answer in place of a raw terminal, in both the MCP and REST
surfaces, so a named lead whose terminal changes can still poll and answer
its own delegations while a different lead is refused both.

Mutation evidence (each one-line change killed a named test, then reverted
to green):
- Principal.ownerKey() PRIMARY case made unconditional (dropped the
  null-name guard) -> ownerKeyCoversEveryRole dies:
  "expected: <null> but was: <leader:null>"
- OBSERVER case changed to use the "worker" prefix -> ownerKeyCoversEveryRole
  dies: "expected: <observer:term_observer> but was: <worker:term_observer>"
- prefixed() changed to drop the role prefix entirely -> both
  ownerKeyCoversEveryRole and rolePrefixesKeepLeadAndArchitectKeysDistinct
  die on a lead/architect key collision: "expected: <leader:opus> but was:
  <opus>"
- PRIMARY case changed to key on terminal instead of name ->
  rolePrefixesKeepLeadAndArchitectKeysDistinct and ownerKeyCoversEveryRole
  die: "expected: <leader:opus> but was: <leader:term_lead>"
- ARCHITECT case changed to key on the slot name instead of terminal ->
  same two tests die: "expected: <architect:opus> but was: <architect:design>"

All five mutations were caught by the existing test suite; no test needed
adding.
2026-10-04 20:42:44 +02:00
Dai Ha aabecce901 Merge remote-tracking branch 'origin/worker/736-presence-forget-f35144-9'
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m58s
2026-10-04 20:19:52 +02:00
Dai Ha 6754b4edbc fleetd #736: release clears the member's presence entry
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 1m45s
SessionManager.releaseRemoved() tore down a member's registry row and pane
but never cleared it from MemberPresence, so a terminal stayed marked
"present" for the daemon's lifetime after release/idle-reap/shutdown drain.
Clear it in the method's unconditional finally block, alongside the other
must-always-run teardown step, so every release path (explicit release,
the idle reaper's releaseIfCurrent, and a shutdown drain) forgets it the
same way, and a throw from the dirty-worktree check does not skip it.

MemberPresence.forget(null) throws NullPointerException (verified empirically:
ConcurrentHashMap.remove(null) NPEs on key.hashCode()), so the new call guards
on a non-null, non-blank terminal id rather than relying on forget to no-op.
2026-10-04 20:15:34 +02:00
Dai Ha 787ae0ed7a fleetd #726: tell the handover skill which rollover behaviour is live
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 1m59s
Unit 2 replaces the /clear continuation with a real process restart, so every
paragraph in the skill that describes /clear goes false the moment the new jar
is deployed. The code is not merged yet, and a merge is not a deployment, so
rewriting those paragraphs now would hand a lead doing a handover tonight a
document that does not match the daemon it is talking to.

Add a dated note instead. It states that the /clear text stays accurate while
the old jar runs, and gives a test a lead can apply with no shell: the
fleet_handover tool description is served by the running daemon, so if it still
says "clear your pane", the old behaviour is live. It also names the two things
that change, including the one that doubles as a second indicator -- the
"never observed as WORKING after 8 consecutive IDLE/DONE polls" warning cannot
appear once the wait that logs it is deleted. The note names the condition for
deleting itself.

Re-measure the roll evidence while here. The skill recorded four
"lead-rollover: rolled" lines from 2026-09-22; the log now holds 20, against a
control of 86 "lead-rollover:" lines, and "Unknown command" still returns 0.
Add the elapsed spread (median 16507 ms, max 48261 ms, two above 45000 ms) with
the caveat that it times the whole roll and is dominated by the wait for the
calling turn to end, so a slow roll is not a failed one.

Markdown only, no code touched, so no build was run.
2026-10-04 20:09:58 +02:00
Dai Ha e3050efe8b fleetd #705: correct the stale reason on the TASK_READ gate
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 50s
CI / build (push) Failing after 1m48s
The comment said ticket ids are a sequential counter with no owner check, so a
holder could walk every ticket and read another session's reply. PRs #712 and
#716 added that owner check: MessageService.ownsTicket compares a ticket's
creatorTerminal to the caller on every read.

The rule is still right, so only the reason changes. This matters now because
fleetd #737 is deciding ticket ownership across a lead handover, and a reader
who believed the old text could delete the TASK_READ restriction on the grounds
that its stated reason no longer applies.

Comment-only. mvn -o clean install: Tests run: 2083, Failures: 0, Errors: 0,
BUILD SUCCESS. Flagged by the #705 option-1 worker as out of its scope, which
was the right call.
2026-10-04 19:59:45 +02:00
Dai Ha 11998cd626 Merge remote-tracking branch 'origin/worker/705-observer-14c258-6'
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 56s
CI / build (push) Failing after 1m54s
2026-10-04 19:52:34 +02:00
Dai Ha 8e5394f63f fleetd #705 option 1: narrow the unconfigured-pane floor to OBSERVER
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Failing after 1m57s
Adds Role.OBSERVER as the bottom rung CallerResolver falls to when a
herdr pane matches no live roster entry, lead, architect slot, or
collaborator tab. An observer may only READ/METRICS and REPLY/ASK on
its own pane. Widens the presence gate so an observer's MCP contact
still marks it deliverable, matching what already happens for a
worker or architect, so a pane that outlives a daemon restart is not
left permanently undeliverable.

Ships as defence in depth alongside the already-merged ticket-owner
check (#712/#716), which closed the reachable exploit this ticket
reported.
2026-10-04 19:46:54 +02:00
Dai Ha 428a12af62 fleetd #722: reconcile presence that arrives before a session's registry entry
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m48s
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 47s
CI / build (push) Failing after 1m46s
A member whose MCP contact lands between launcher.spawn() and registry.put()
had its presence marked, but the SPAWNING -> READY transition that markPresent
triggers found no registry entry yet and silently did nothing. The mark then
persisted while registration left the session in SPAWNING, with nothing to
retry the transition. That left the session undeliverable to reclaim/seat
accounting even though it was present and deliverable.

Add SessionManager.reconcilePresence, called right after registry.put in both
the plain-spawn and worktree-spawn paths, to retry the transition for a
terminal already marked present. One private helper serves both call sites.

Tests cover both orderings (contact-then-register and register-then-contact)
for both spawn paths, plus a terminal never marked present staying in
SPAWNING. The contact-then-register tests use a new PresenceRacingLauncher
test double that marks presence from inside spawn(), before acquire()'s own
registry.put runs.
2026-10-04 19:23:02 +02:00
Dai Ha a332dfdb2c Merge remote-tracking branch 'origin/worker/726-ea34a0-2'
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m53s
2026-10-04 19:09:03 +02:00
Dai Ha d0f4ae057b fleetd #726 unit 3 fix: release the single-flight claim when continuationRunner rejects the hand-off
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 1m51s
confirm() takes the per-lead-terminal claim before handing the roll to
continuationRunner, and the only release path was runRollover's own
finally. If continuationRunner.accept itself throws, runRollover never
starts, so that finally never runs, and nothing else ever writes
rollingByTerminal — the claim is held forever and the terminal can never
be rolled again. This differs from fleetd #615, which covers a throw
INSIDE the continuation (runRollover already catches that and still
releases the claim) — this is a throw from the hand-off itself, which
is not reachable with today's virtual-thread runner but would be with
a bounded executor's RejectedExecutionException.

confirm() now catches that throw, releases the claim, and overwrites
the IN_PROGRESS outcome with a terminal FAILED one, matching how a
throw inside the continuation is already surfaced.
2026-10-04 19:06:57 +02:00
Dai Ha 7b3beaa209 Merge remote-tracking branch 'origin/worker/726-10cbf0-1'
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m53s
2026-10-04 19:06:47 +02:00
Dai Ha b3b2bf3da6 fleetd #726 unit 1 review fixes: correct relaunch's javadoc and dedupe its resolve logic
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Failing after 1m49s
relaunch's javadoc said it reads the live config; it actually reads the
FleetConfig snapshot this launcher was constructed with (fleet.leaders
is the frozen half), so say that and note a profile/tab edit needs a
daemon restart.

relaunch's recognise-only refusal reused ensureLeads()'s log wording,
which claims the lead 'is not live' — true in ensureLeads()'s context
(reached only after a short live count), false in relaunch's (which
never counts liveness, by design). Dropped that clause.

Pulled the declared/creatable/profile-configured resolution shared by
ensureLeads() and relaunch() into one private resolveLaunchable(name)
helper (returns a new ResolvedLead(lead, profile) record, or null
having logged), so the three refusals and their wording live in one
place instead of two copies that can drift. Behaviour-preserving:
ensureLeads() keeps its own liveness-count logic around the shared
resolve, and the existing 37 LeadLauncherTest cases are unchanged and
still pass.
2026-10-04 18:59:56 +02:00
Dai Ha 8bb2aa0be4 Merge remote-tracking branch 'origin/worker/729-5961c6-3'
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 2m0s
2026-10-04 18:59:54 +02:00
Dai Ha 4ce3149bfd fleetd #726 unit 3: make LeadRollover.confirm() single-flight per lead terminal
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m48s
Two open() calls for the same lead terminal minted two tokens that both
passed confirm()'s ownership check, so both could reach the deferred
continuation and roll the same lead twice. confirm() now claims a
per-lead-terminal slot (an atomic put-if-absent) once every other gate has
passed, refusing a concurrent confirm with the new ROLL_ALREADY_RUNNING
reason; runRollover releases the claim in a finally, on both the success
and the thrown-exception path.
2026-10-04 18:53:53 +02:00
Dai Ha 1fc9e85bf1 fleetd #729: fold a per-boot nonce into every turnId
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 1m57s
askSeq restarts at 0 on every daemon boot, so a turnId (session#n)
minted by one Rendezvous instance could be minted again by a later
instance and resolve to an unrelated ask. Fold a per-instance nonce
into the mint, the same way #719 fixed MessageService's ticket ids.
2026-10-04 18:53:05 +02:00
Dai Ha 38544d467c fleetd #726 unit 1: give LeadLauncher a public single-lead relaunch seam
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Failing after 1m48s
Adds LeadLauncher.relaunch(name), which starts exactly the named lead
from the live config, outside of ensureLeads()'s instances bookkeeping.
It retries the whole launch attempt (not just the agent_name_taken/
agent_pane_busy cases ResilientAgentLaunch already retries inside one
agents.start call) up to RELAUNCH_ATTEMPTS times.

launch() now returns the started Agent (null on failure) instead of a
boolean, so relaunch() and ensureLeads() share the same primitive.
2026-10-04 18:52:38 +02:00
Dai Ha 809b7d9b20 Merge remote-tracking branch 'origin/worker/727-ee14ed-3'
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 47s
CI / build (push) Failing after 1m56s
2026-10-04 18:29:28 +02:00
Dai Ha 8cf7215d56 Merge remote-tracking branch 'origin/worker/719-bdd95e-4'
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 1m58s
2026-10-04 18:19:55 +02:00
Dai Ha cf0c9b9316 fleetd #719: make the foreign-id test reach a colliding sequence number
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 1m49s
other.sendAsync had never been called, so other's tasks map was empty and
poll(ticket) returned null regardless of whether the nonce existed — the
test passed against an empty map, not against a colliding id. Mint once on
other so it reaches the same sequence number as the first instance, making
the test exercise the actual collision the nonce guards against.
2026-10-04 18:14:01 +02:00
Dai Ha 337dbd491e fleetd #719: fold a per-boot nonce into every ticket id
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 2m0s
ticketSeq restarted at zero on every daemon boot with no persistence, so a
ticket id minted in one boot could be reused by a later boot and resolve to
an unrelated Task instead of failing to resolve at all. Mint each
MessageService instance's own short nonce once and fold it into every ticket
(task-<nonce>-<n>), so an id from one instance can never match another's id
space.

Adds a disjoint-id-space test and a foreign-instance-ticket test (with the
positive control) in MessageServiceTest.
2026-10-04 18:05:06 +02:00
38 changed files with 2554 additions and 1199 deletions
+35 -7
View File
@@ -159,6 +159,29 @@ fails.
`{action: "cancel", token}` drops a pending request without rolling.
**A change is coming: the roll will restart the process instead of sending `/clear` (fleetd #726
unit 2, written 2026-10-04).**
Today a roll types `/clear` into your pane. Your `claude` process keeps running, so a newer CLI on
disk is never loaded. Unit 2 replaces that: the daemon ends the old pane, launches a fresh one,
waits for the new terminal to be recognised as a lead, and only then sends the bootstrap text.
**Everything below about `/clear` is accurate while the old jar is running.** Unit 2 was not merged
when this note was written, and a merge is not a deployment.
**How to tell which one is live: read your own tool list.** If `fleet_handover`'s description says
it will "clear your pane", the daemon is serving the old behaviour. If it names a restart, the new
behaviour is live. The description comes from the running daemon, so it cannot disagree with the
code that is actually loaded.
Two things change for you once it is live. The `status` outcomes are different: three new failures
replace the `/clear` ones. And the "never observed as WORKING after 8 consecutive IDLE/DONE polls"
warning described below can no longer appear, because that wait is deleted — so if you still see
it, the old jar is running. `TURN_NEVER_SETTLED` does not change, and still means nothing was
touched.
**Delete this note and rewrite the `/clear` paragraphs once the new jar is live.**
**Things that will surprise you:**
- **`accepted` does not mean your pane has been cleared.** It means every gate passed and the roll
@@ -200,13 +223,18 @@ fails.
- **The roll can still refuse after `confirm` returns**, and by then there is no caller to tell.
Those outcomes are logged only, as `lead-rollover:` lines in the daemon log.
- **The bootstrap prompt works end to end. Measured 2026-09-22.** This used to say the fix was
unproven (fleetd #489) and told you to expect a failure. That is no longer true. The daemon log
now holds four `lead-rollover: rolled` lines, and three of them ran on 2026-09-22 at 10:01:43,
10:38:28 and 11:15:47. Each one cleared the old lead and started a fresh session against the
handover file, with the configured `bootstrapText` arriving as its first message. No context was
lost. The old `Unknown command: /clearFresh` failure from 2026-09-12 does not appear in the log
at all. Re-measure both numbers with:
- **The bootstrap prompt works end to end. Measured 2026-09-22, re-measured 2026-10-04.** This used
to say the fix was unproven (fleetd #489) and told you to expect a failure. That is no longer
true. On 2026-09-22 the daemon log held four `lead-rollover: rolled` lines. On 2026-10-04 it holds
**20**, against a control of 86 `lead-rollover:` lines. Each roll cleared the old lead and started
a fresh session against the handover file, with the configured `bootstrapText` arriving as its
first message. No context was lost. The old `Unknown command: /clearFresh` failure from 2026-09-12
does not appear in the log at all.
19 of the 20 carry an `elapsedMs`: median 16507 ms, maximum 48261 ms, and two above 45000 ms. That
figure times the **whole** roll, and the wait for your own turn to end dominates it, so do not
read it as the cost of the clear. Expect a roll to take tens of seconds, and do not treat a slow
one as a failed one. Re-measure all of these with:
```bash
grep -c "lead-rollover: rolled" fleetd/fleetd.out # successful rolls
+19 -14
View File
@@ -27,26 +27,31 @@ through its `fleet_*` tools. No session addresses a peer, a broker, or the netwo
**Every role reads this file.** A member runs in a git worktree of this same repo, so it inherits
this `CLAUDE.md` verbatim, and every rule below is role-conditional.
**Call `fleet_whoami`.** It returns `primary`, `worker`, `architect`, or `collaborator`, resolved by
the daemon from your connection — unforgeable, and the same resolution its authorization gate uses.
A worker also carries its `sessionId`, `profile`, `worktree` and `branch`; an architect carries the
slot name it was bound to; a collaborator carries its registry name and its own `sessionId`, and
**no `leader` key** — a collaborator is a named peer, not a primary. Don't infer what you can ask.
**Call `fleet_whoami`.** It returns `primary`, `worker`, `architect`, `collaborator`, or `observer`,
resolved by the daemon from your connection — unforgeable, and the same resolution its authorization
gate uses. A worker also carries its `sessionId`, `profile`, `worktree` and `branch`; an architect
carries the slot name it was bound to; a collaborator carries its registry name and its own
`sessionId`, and **no `leader` key** — a collaborator is a named peer, not a primary. An **observer**
carries only its own `sessionId`: a pane the daemon could not place as any of the above, authorized
to `READ`/`METRICS` and to `REPLY`/`ASK` on its own pane and nothing more — never `SEND`, never a
ticket. Don't infer what you can ask.
Only if that call is unavailable, fall back to these — each is one-way, so keep reading until one
fires: the reply charter in your system prompt (*"You are a spawned member in the
claude-bridge fleet"*) ⇒ **spawned member**; fleet tools prefixed `mcp__fleet__*` ⇒ **spawned
member** (the launcher fixes that mount name; a primary's mount is named by whoever wrote its
`.mcp.json`, so it varies — and a member spawned before CB-632 still says `mcp__bridge__*`); `ANTHROPIC_BASE_URL` set ⇒ **spawned member** (Claude-model members run
on a clean env, so its *absence* proves nothing). None of these separate a worker from an architect —
only `fleet_whoami` does. **And none of them fires for a collaborator at all**: every signal in the
ladder detects a *spawned* member, while a collaborator is a tab a person opened by hand, so it has
no charter, no fixed mount name and a normal environment. A collaborator that cannot call
`fleet_whoami` therefore falls to the line below and acts as a worker. That is the safe direction —
it under-privileges, and the refusals are loud — but it means a collaborator has no way to learn
what it is except by asking. **Still unsure ⇒ act as a worker**, the most restricted member role. The
two mistakes are not symmetric: a primary acting as a worker is refused by the authorization gate —
loud and self-correcting — while a member acting as the primary ends its turn with no `fleet_reply`,
on a clean env, so its *absence* proves nothing). None of these separate a worker from an architect,
or a worker from an **observer** — an observer is just as unspawned as a collaborator and carries
none of these signals either, so only `fleet_whoami` tells the two apart. **And none of them fires
for a collaborator at all**: every signal in the ladder detects a *spawned* member, while a
collaborator is a tab a person opened by hand, so it has no charter, no fixed mount name and a
normal environment. A collaborator — or an observer — that cannot call `fleet_whoami` therefore falls
to the line below and acts as a worker. That is the safe direction — it under-privileges, and the
refusals are loud — but it means a collaborator or an observer has no way to learn what it is except
by asking. **Still unsure ⇒ act as a worker**, the most restricted member role this ladder can name.
The two mistakes are not symmetric: a primary acting as a worker is refused by the authorization gate
— loud and self-correcting — while a member acting as the primary ends its turn with no `fleet_reply`,
and the sender silently receives nothing. Fail toward the recoverable error.
### Invariants — every role, no exceptions
+22 -17
View File
@@ -99,17 +99,17 @@ bind:
# contextHighNudge: false
# Lead rollover (fleetd #480): replace a lead session that has decided it is ready to be replaced,
# without an operator doing it by hand. A lead writes a handover file, then asks fleetd to clear its
# own pane and bootstrap a fresh session against that file.
# without an operator doing it by hand. A lead writes a handover file, then asks fleetd to end its
# own pane, launch a fresh one, and bootstrap that fresh session against the handover file.
#
# Opt-in on purpose — it clears the lead's own pane on request, so upgrading the daemon must never
# acquire that ability for you. Absent block = feature off, and nothing is constructed at all. Even
# once present, nothing but an explicit confirm() call — one that passes every check — can ever
# cause a /clear: there is no recurring timer, heartbeat or scheduler anywhere in this feature that
# fires one on its own initiative. confirm() itself is called FROM the calling lead's own turn, so
# it cannot clear the pane inline (that pane is still WORKING); instead it schedules a one-shot
# Opt-in on purpose — it tears down the lead's own pane on request, so upgrading the daemon must
# never acquire that ability for you. Absent block = feature off, and nothing is constructed at all.
# Even once present, nothing but an explicit confirm() call — one that passes every check — can ever
# tear a pane down: there is no recurring timer, heartbeat or scheduler anywhere in this feature that
# fires one on its own initiative. confirm() itself is called FROM the calling lead's own turn, so it
# cannot act on the pane inline (that pane is still WORKING); instead it schedules a one-shot
# continuation that waits for the SAME confirm() call's turn to end, then does the actual work. See
# dev.ltms.fleet.lead.LeadRollover's class javadoc for the exact order (fleetd #480 correction).
# dev.ltms.fleet.lead.LeadRollover's class javadoc for the exact order.
#
# handoverPath: REQUIRED when this block is present — where the handover file a fresh lead session
# reads must live. No default (an operator-specific path); a present block with no
@@ -118,25 +118,30 @@ bind:
# working directory when that lead has none configured) — never against whatever
# directory the daemon process happens to have been started in. An absolute path is
# used unchanged. Prefer an absolute path if the daemon and the lead's pane might not
# share a working directory (fleetd #480 follow-up).
# share a working directory.
# requireOperatorConfirm: true # default true — confirm() refuses unless the caller also passes
# # operatorConfirmed: true
# maxDocAgeSeconds: 3600 # default 3600 — refuse a handover file older than this
# turnSettleSeconds: 20 # default 20 — how long the deferred roll waits for the CALLING
# # lead's own turn to end (its pane to report injectable again)
# # before sending /clear at all. If this elapses, /clear is NEVER
# # sent — a lead that never goes idle is still doing real work.
# clearSettleSeconds: 20 # default 20 — how long to wait for the pane to become injectable
# # again AFTER /clear before giving up (never sends bootstrapText
# # if this elapses). A separate, second wait from turnSettleSeconds.
# # before tearing the old pane down at all. If this elapses, nothing
# # is torn down — a lead that never goes idle is still doing real
# # work.
# relaunchReadySeconds: 45 # default 45 — bounds two later waits, after the old pane is gone
# # and a fresh one has been launched: first, for the fresh pane to
# # reach a real turn boundary (never sends bootstrapText if THIS one
# # elapses); second, for the new terminal to be recognised as this
# # lead (bootstrapText is sent either way once the first wait
# # passes). A separate, later pair of waits from turnSettleSeconds.
# bootstrapText: "..." # default names the RESOLVED (absolute) handoverPath — sent to
# # the lead once its pane settles after /clear
# # the freshly relaunched lead's pane once it reaches a real turn
# # boundary
# leadRollover:
# handoverPath: /path/to/handover.md
# requireOperatorConfirm: true
# maxDocAgeSeconds: 3600
# turnSettleSeconds: 20
# clearSettleSeconds: 20
# relaunchReadySeconds: 45
# bootstrapText: "Fresh lead session: read the handover file and carry on."
# Fleet health detection is dormant unless enabled (CB-573). It reads one whole-fleet agent list
@@ -226,9 +226,8 @@ public final class Fleetd {
* is the same way — a person's own tab, matched to a configured name, never spawned.
*
* <p>Neither a lead nor a collaborator is ever enrolled in {@link MemberPresence} — {@code
* FleetMcp} marks presence for every spawned member (worker and architect), deliberately, since
* that map doubles as the member roster's availability signal and a lead or collaborator counted
* there would show up as an available member. So without the second and third disjuncts a lead or
* FleetMcp} marks presence only for a worker, an architect, or the unconfigured-pane floor,
* never for a lead or a collaborator. So without the second and third disjuncts a lead or
* collaborator is permanently un-deliverable: every send to one sat on the gate for
* {@code READINESS_GRACE_POLLS} (~60s) and then failed having never been typed into the pane.
*
@@ -942,21 +941,27 @@ public final class Fleetd {
* cfg.leadHeartbeat()}
* @param leadAgents the {@link AgentControl} instance that reaches the LEAD's pane (not
* {@code memberAgents}), normally {@code router.leadAgents()}
* @param leadSpaces the {@link WorkspaceControl} instance that reaches the LEAD's
* workspace, normally {@code router.leadSpaces()} — used to tear down
* a rolled lead's old pane and confirm it is gone
* @param launcher starts the fresh lead a roll relaunches once the old one is gone
* @param config the live {@link ConfigRef}, captured only inside the returned
* supplier and the workspace lookup — never dereferenced here
* supplier and the two lookups below — never dereferenced here
* @param liveLeadTerminals terminal id → lead NAME for every CURRENTLY recognised lead, normally
* the same {@code leads} supplier {@code main} already builds for
* {@code HerdrRouter}/{@link #leadSeatLookup} — never a value snapshot
* @return a constructed {@link LeadRollover}, or {@code null} when {@code leadRollover:} is
* absent from the startup config
*/
static LeadRollover leadRollover(FleetConfig cfg, AgentControl leadAgents, ConfigRef config,
static LeadRollover leadRollover(FleetConfig cfg, AgentControl leadAgents,
WorkspaceControl leadSpaces, LeadLauncher launcher, ConfigRef config,
Supplier<Map<String, String>> liveLeadTerminals) {
if (cfg.leadRollover() == null) {
return null;
}
Function<String, String> leadNameForTerminal = terminal -> liveLeadTerminals.get().get(terminal);
Function<String, String> leadWorkspace = terminal -> {
String leadName = liveLeadTerminals.get().get(terminal);
String leadName = leadNameForTerminal.apply(terminal);
if (leadName == null) {
return null;
}
@@ -964,7 +969,8 @@ public final class Fleetd {
FleetConfig.Leader leader = fleet == null ? null : fleet.leaders().get(leadName);
return leader == null ? null : leader.cwd();
};
return new LeadRollover(leadAgents, () -> config.get().leadRollover(), leadWorkspace);
return new LeadRollover(leadAgents, leadSpaces, launcher, () -> config.get().leadRollover(),
leadWorkspace, leadNameForTerminal, liveLeadTerminals);
}
/**
@@ -297,11 +297,15 @@ final class FleetdAssembly {
leadsRef.set(leads);
collaboratorTerminalsRef.set(collaboratorTerminals);
// Constructed unconditionally — it is cheap and side-effect free — so a LeadRollover built
// below can relaunch a lead even on a boot where herdr was down for the ensureLeads() call.
LeadLauncher leadLauncher = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg);
// CB-558: start any declared lead that is not already running. After the scanner is built,
// and only when herdr answered — the launcher's whole safety property is that it can count
// live leads first, and must never guess and risk a second orchestrator.
if (herdrUp && !leaders.isEmpty()) {
int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads();
int launched = leadLauncher.ensureLeads();
if (launched > 0) {
log.info("lead auto-launch: {} lead(s) started", launched);
}
@@ -440,7 +444,8 @@ final class FleetdAssembly {
heartbeatScheduler.shutdownNow();
}
// fleetd #480: lead rollover. Opt-in; absent `leadRollover:` this is never constructed.
LeadRollover leadRollover = Fleetd.leadRollover(cfg, router.leadAgents(), config, leads);
LeadRollover leadRollover = Fleetd.leadRollover(cfg, router.leadAgents(), router.leadSpaces(),
leadLauncher, config, leads);
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
pushLoop, metrics);
@@ -138,13 +138,15 @@ public final class Authz {
// fleet_whoami — and carries no secrets: no ticket reply, no pending question, and no
// other session's turn state. Those live under TASK_READ. METRICS is the separate
// Prometheus scrape. Both are open to every authenticated role, including a
// collaborator.
// collaborator and the unconfigured-pane floor.
case READ, METRICS -> caller.isPrimary() || caller.isWorker() || caller.isArchitect()
|| caller.isCollaborator();
|| caller.isCollaborator() || caller.isObserver();
// Ticket polling and session status, open to every role READ is open to except a
// collaborator: ticket ids are a sequential counter with no owner check, so a holder
// could walk every ticket and read another session's delegation reply.
// collaborator or an observer. MessageService compares a ticket's creator to the
// caller on every read as well, so dropping this gate would not expose another
// session's reply — it would move the refusal later and widen what a caller that
// never orchestrates can probe.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect
@@ -40,7 +40,7 @@ import java.util.function.Supplier;
* the case the previous step does not catch: a binding with no live spawned-member session.</li>
* <li>A loopback peer PID that maps to an operator-labelled collaborator tab ⇒
* {@link Role#COLLABORATOR}, carrying that collaborator's name.</li>
* <li>A loopback peer PID that maps to any other herdr pane ⇒ {@link Role#WORKER}. This is
* <li>A loopback peer PID that maps to any other herdr pane ⇒ {@link Role#OBSERVER}. This is
* unforgeable (the OS reports the PID, herdr owns the PID→pane map) and is honoured
* regardless of auth mode, so enabling auth never breaks the fleet.</li>
* <li>Otherwise, under {@code token} mode, a valid bearer token ⇒ {@link Role#PRIMARY}.</li>
@@ -323,7 +323,7 @@ public final class CallerResolver {
// of the above keeps that stronger role.
return Principal.collaborator(collaborator, c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
return Principal.observer(c.terminal(), c.pid()); // unforgeable; never token-gated
}
if (tokenMode) {
@@ -88,6 +88,13 @@ public record Principal(Role role, String terminal, long pid, String name) {
return new Principal(Role.COLLABORATOR, terminal, pid, name);
}
/**
* The unconfigured-pane floor: a loopback caller whose pane matched no other role.
*/
public static Principal observer(String terminal, long pid) {
return new Principal(Role.OBSERVER, terminal, pid);
}
public boolean isPrimary() {
return role == Role.PRIMARY;
}
@@ -104,11 +111,15 @@ public record Principal(Role role, String terminal, long pid, String name) {
return role == Role.WORKER;
}
public boolean isObserver() {
return role == Role.OBSERVER;
}
/**
* Whether this caller is a spawned member with its own pane.
*
* <p>Both workers and architects are spawned members. A lead is excluded because recording it
* as present would count it as an available member in the roster.
* <p>Both workers and architects are spawned members. A lead is not: it is a peer the
* operator started and named, never a pane this daemon spawned.
*/
public boolean isSpawnedMember() {
return role == Role.WORKER || role == Role.ARCHITECT;
@@ -134,12 +145,32 @@ public record Principal(Role role, String terminal, long pid, String name) {
return terminal != null && terminal.equals(sessionId);
}
/**
* Stable identity used to own tickets and open turns. The unnamed primary has no owner key so
* it can use the message layer's primary-wide ticket access rule.
*/
public String ownerKey() {
return switch (role) {
case PRIMARY -> name == null ? null : prefixed("leader", name);
case WORKER -> prefixed("worker", terminal);
case ARCHITECT -> prefixed("architect", terminal);
case COLLABORATOR -> prefixed("collaborator", name);
case OBSERVER -> prefixed("observer", terminal);
case ANONYMOUS -> "anonymous";
};
}
private static String prefixed(String role, String identity) {
return role + ":" + identity;
}
/** Short, non-sensitive description for audit lines and error details. */
public String describe() {
return switch (role) {
case WORKER -> "worker:" + terminal;
case ARCHITECT -> "architect:" + name;
case COLLABORATOR -> "collaborator:" + name;
case OBSERVER -> "observer:" + terminal;
case PRIMARY -> name == null ? "primary" : "leader:" + name;
case ANONYMOUS -> "anonymous";
};
@@ -45,6 +45,17 @@ public enum Role {
*/
COLLABORATOR,
/**
* A loopback pane that resolved to none of the roles above: not a live spawned member, not a
* configured lead, not a bound architect slot, not a configured collaborator tab. Unforgeable
* like a worker's — derived from the connection's pane, never from a request argument, and
* honoured regardless of auth mode. May {@code READ} and {@code METRICS}, and {@code REPLY}/
* {@code ASK} only as its own pane; may not {@code SPAWN}/{@code STOP}/{@code DRAIN}/
* {@code HANDOVER}, {@code SEND}, poll a ticket ({@code TASK_READ}), or reach the coordination
* broker ({@code COORD_SEND}/{@code COORD_READ}).
*/
OBSERVER,
/** Authenticated as nothing. Authorized for nothing but {@code /healthz}. */
ANONYMOUS
}
@@ -67,7 +67,7 @@ import java.util.function.Supplier;
* {@code models:} above: {@code dev.ltms.fleet.lead.LeadRollover} holds a
* {@code Supplier<FleetConfig.LeadRollover>} (the same {@code () -> config.get().x()} shape)
* and reads {@code handoverPath}/{@code requireOperatorConfirm}/{@code maxDocAgeSeconds}/
* {@code turnSettleSeconds}/{@code clearSettleSeconds}/{@code bootstrapText} fresh on every
* {@code turnSettleSeconds}/{@code relaunchReadySeconds}/{@code bootstrapText} fresh on every
* {@code open()}/{@code confirm()} call (and on the deferred post-{@code confirm()}
* continuation fleetd #480's correction added — see {@code LeadRollover}'s class doc) rather
* than capturing them into fields at construction — unlike its closest
@@ -1416,17 +1416,16 @@ public record FleetConfig(
* before anything exists to call — the same fact already true of adding a brand-new
* {@code profiles:} entry.
*
* <p><strong>{@code turnSettleSeconds} (fleetd #480 correction):</strong> {@code confirm()} is
* called FROM the calling lead's own turn, so its pane is still {@code WORKING} the instant
* {@code confirm()} validates every gate and schedules the roll. {@code
* dev.ltms.fleet.lead.LeadRollover}'s deferred continuation waits up to this many seconds for
* that SAME pane to report {@code IDLE} or {@code DONE} — i.e. for the calling turn to actually
* end — before it sends {@code /clear} at all. {@code BLOCKED} does not count: that is a live
* turn merely paused, not one that has finished. If that wait times out, no {@code /clear} is
* ever sent: a lead that never goes idle is still doing real work, and clearing it would
* destroy live context. This is a separate wait from {@code clearSettleSeconds} below, which
* bounds the SECOND wait, for the pane to reach {@code IDLE} or {@code DONE} again AFTER
* {@code /clear} has already gone out.
* <p><strong>{@code turnSettleSeconds}:</strong> {@code confirm()} is called FROM the calling
* lead's own turn, so its pane is still {@code WORKING} the instant {@code confirm()} validates
* every gate and schedules the roll. {@code dev.ltms.fleet.lead.LeadRollover}'s deferred
* continuation waits up to this many seconds for that SAME pane to report {@code IDLE} or
* {@code DONE} — i.e. for the calling turn to actually end — before it ends the old pane's
* process at all. {@code BLOCKED} does not count: that is a live turn merely paused, not one
* that has finished. If that wait times out, the old pane is never touched: a lead that never
* goes idle is still doing real work, and the roll ends that pane's whole process — there is no
* way back from this once it runs, so this wait is the only thing standing between "still
* working" and "gone".
*
* @param handoverPath required when this block is present — where the handover file a fresh
* lead session reads must live. There is no sane non-null default for an
@@ -1447,36 +1446,49 @@ public record FleetConfig(
* attempt can never be mistaken for a fresh one.
* @param turnSettleSeconds default 300 — bound on how long the deferred roll waits for the
* CALLING lead's own turn to end (its pane to report {@code IDLE} or
* {@code DONE}) before sending {@code /clear} at all. See the paragraph
* above.
* @param clearSettleSeconds default 20 — bound on how long to wait for the lead's pane to
* report {@code IDLE} or {@code DONE} again after {@code /clear} before
* giving up. A roll that times out here never sends {@code bootstrapText}.
* {@code DONE}) before ending that pane's process at all. See the
* paragraph above.
* @param relaunchReadySeconds default 45 — bound on EACH of two separate waits that run after
* the old lead's pane has been torn down and a fresh one launched: first,
* for the fresh pane itself to reach a real turn boundary ({@code IDLE} or
* {@code DONE}, never merely {@code BLOCKED}) — the safety gate, since
* typing into a pane that has not finished booting loses the keystrokes;
* second, for the fresh terminal to show up as a recognised lead, which is
* bookkeeping rather than a safety gate, so a timeout on this second wait
* does not withhold {@code bootstrapText} — it is sent once the pane is
* ready regardless. Recognition comes from the same periodically-refreshed
* scan {@code LeadTabScanner} already keeps ({@code scanIntervalSeconds},
* 10s live), so a budget has to clear more than one scan interval to leave
* any real margin for the CLI's own boot time; 20 was rejected for exactly
* that reason — at a 10s scan interval it only buys two scans. 45 buys
* roughly four. Only a timeout on the FIRST wait (the pane never becomes
* ready) withholds {@code bootstrapText}.
* @param bootstrapText default a sentence naming the RESOLVED handover path — sent to the
* lead's pane once it settles after {@code /clear}, telling the fresh
* session where to read the handover and carry on. Left {@code null} here
* when the operator configures none: the default sentence cannot be built
* at construction time because it must name the path AFTER {@code
* dev.ltms.fleet.lead.LeadRollover#open} has resolved a relative {@code
* handoverPath} against the calling lead's workspace, which this record has
* no way to know — see {@link #bootstrapTextFor(String)}.
* fresh lead's pane once it reaches a real turn boundary after relaunch,
* telling the fresh session where to read the handover and carry on. Left
* {@code null} here when the operator configures none: the default sentence
* cannot be built at construction time because it must name the path AFTER
* {@code dev.ltms.fleet.lead.LeadRollover#open} has resolved a relative
* {@code handoverPath} against the calling lead's workspace, which this
* record has no way to know — see {@link #bootstrapTextFor(String)}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record LeadRollover(String handoverPath, Boolean requireOperatorConfirm,
Integer maxDocAgeSeconds, Integer turnSettleSeconds,
Integer clearSettleSeconds, String bootstrapText) {
Integer relaunchReadySeconds, String bootstrapText) {
public LeadRollover {
requireOperatorConfirm = requireOperatorConfirm == null || requireOperatorConfirm;
maxDocAgeSeconds = (maxDocAgeSeconds == null || maxDocAgeSeconds <= 0) ? 3600 : maxDocAgeSeconds;
turnSettleSeconds = (turnSettleSeconds == null || turnSettleSeconds <= 0) ? 300 : turnSettleSeconds;
clearSettleSeconds = (clearSettleSeconds == null || clearSettleSeconds <= 0) ? 20 : clearSettleSeconds;
relaunchReadySeconds = (relaunchReadySeconds == null || relaunchReadySeconds <= 0)
? 45 : relaunchReadySeconds;
bootstrapText = (bootstrapText == null || bootstrapText.isBlank()) ? null : bootstrapText;
}
/**
* The text actually sent to the lead's pane once it settles after {@code /clear}: the
* operator's configured {@link #bootstrapText} when one is set, otherwise the default
* sentence built from {@code resolvedHandoverPath}.
* The text actually sent to the fresh lead's pane once it reaches a real turn boundary
* after relaunch: the operator's configured {@link #bootstrapText} when one is set,
* otherwise the default sentence built from {@code resolvedHandoverPath}.
*
* @param resolvedHandoverPath the ABSOLUTE path {@code dev.ltms.fleet.lead.LeadRollover
* #open} already resolved — never the raw configured {@link
@@ -1919,6 +1931,7 @@ public record FleetConfig(
rejectNegativeMaxLoad(yaml);
rejectAutoCompactWindowOutOfRange(yaml);
warnConflictingAutoCompactWindows(yaml);
warnRetiredClearSettleSecondsKey(yaml);
rejectMalformedProfilePatterns(yaml);
rejectUnknownKind(yaml);
rejectUnknownAuthMode(yaml);
@@ -2342,6 +2355,33 @@ public record FleetConfig(
names, String.join(", ", detail));
}
/**
* Warn when a {@code leadRollover:} block still sets the retired {@code clearSettleSeconds}
* key. {@link LeadRollover} carries {@code @JsonIgnoreProperties(ignoreUnknown = true)} and no
* longer declares that component, so Jackson drops it with no signal of its own — this raw-YAML
* check is the only place an operator's now-inert setting is reported at all; by the time a
* {@link LeadRollover} instance exists to run a validator against, the key is already gone.
*
* @param yaml the raw config text
*/
static void warnRetiredClearSettleSecondsKey(String yaml) {
Map<?, ?> raw;
try {
raw = YAML.readValue(yaml, Map.class);
} catch (IOException | IllegalArgumentException e) {
return;
}
if (raw == null || !(raw.get("leadRollover") instanceof Map<?, ?> leadRollover)) {
return;
}
if (leadRollover.containsKey("clearSettleSeconds")) {
log.warn("leadRollover.clearSettleSeconds is retired and no longer read. Set "
+ "leadRollover.relaunchReadySeconds instead: it bounds how long to wait, after "
+ "a lead is relaunched, for its pane to become ready and then for it to be "
+ "recognised as a lead. Remove clearSettleSeconds from fleetd.yaml.");
}
}
/**
* Reject a profile whose {@code errorPattern} (fleetd #201 Unit 5) or {@code exhaustedPattern}
* (CB-578 stage A) is not a valid Java regex, naming the profile, the key, and the parser's own
@@ -83,6 +83,9 @@ public final class LeadLauncher {
private static final Logger log = LoggerFactory.getLogger(LeadLauncher.class);
/** Attempts {@link #relaunch(String)} makes before giving up and returning {@code null}. */
static final int RELAUNCH_ATTEMPTS = 3;
private final AgentControl agents;
private final WorkspaceControl spaces;
private final FleetConfig cfg;
@@ -201,24 +204,14 @@ public final class LeadLauncher {
log.info("lead '{}': {} live, {} wanted — nothing to start", name, running, wanted);
continue;
}
if (!lead.isCreatable()) {
// A lead with a `tab:` but no `profile:` is recognise-only by design: the operator
// opens it by hand. Say so once rather than looking like a silent failure.
log.info("lead '{}' is not live, and names no profile — it can be recognised but not "
+ "launched. Add `profile:` under fleet.leaders.{} to have fleetd start it.",
name, name);
continue;
}
FleetConfig.Profile profile = cfg.profiles().get(lead.profile());
if (profile == null) {
log.warn("lead '{}' names profile '{}', which is not configured — not launching",
name, lead.profile());
ResolvedLead resolved = resolveLaunchable(name);
if (resolved == null) {
continue;
}
for (int i = running; i < wanted; i++) {
if (launch(name, lead, profile)) {
if (launch(name, resolved.lead(), resolved.profile()) != null) {
started++;
}
}
@@ -226,6 +219,82 @@ public final class LeadLauncher {
return started;
}
/** A declared lead paired with the profile it launches on — {@link #resolveLaunchable}'s result. */
private record ResolvedLead(FleetConfig.Leader lead, FleetConfig.Profile profile) {
}
/**
* The declared {@code Leader} and its {@code Profile} for {@code name}, read from the config
* snapshot this launcher was constructed with.
*
* @return the resolved pair, or {@code null} (having logged) if {@code name} is not declared
* under {@code fleet.leaders}, that lead names no {@code profile:} (a {@code tab:}-only,
* recognise-only lead), or its {@code profile:} is not configured. Shared by
* {@link #ensureLeads()} and {@link #relaunch(String)} so the three refusals and their
* wording live in one place.
*/
private ResolvedLead resolveLaunchable(String name) {
FleetConfig.Leader lead = cfg.fleet().leaders().get(name);
if (lead == null) {
log.warn("lead '{}' is not declared under fleet.leaders — not launching", name);
return null;
}
if (!lead.isCreatable()) {
// A lead with a `tab:` but no `profile:` is recognise-only by design: the operator
// opens it by hand. Say so once rather than looking like a silent failure.
log.info("lead '{}' names no profile — it can be recognised but not launched. Add "
+ "`profile:` under fleet.leaders.{} to have fleetd start it.", name, name);
return null;
}
FleetConfig.Profile profile = cfg.profiles().get(lead.profile());
if (profile == null) {
log.warn("lead '{}' names profile '{}', which is not configured — not launching",
name, lead.profile());
return null;
}
return new ResolvedLead(lead, profile);
}
/**
* Start the named lead from the config snapshot this launcher was constructed with — not a
* live read, so a lead's {@code profile:} or {@code tab:} edited in config needs a daemon
* restart to take effect here — outside of {@link #ensureLeads()}'s {@code instances}
* bookkeeping.
*
* @return the started {@link Agent}, or {@code null} if {@code name} is not declared under
* {@code fleet.leaders}, that lead names no {@code profile:} (a {@code tab:}-only,
* recognise-only lead), its {@code profile:} is not configured, or every attempt up to
* {@link #RELAUNCH_ATTEMPTS} failed to start it. Never throws.
*
* <p>Does not count how many instances of this lead are already live. {@link #ensureLeads()}'s
* count exists to avoid starting a second orchestrator; the caller of this method has already
* decided to replace the lead and owns that decision.
*
* <p>Retries the whole launch attempt — not only the {@code agent_name_taken}/
* {@code agent_pane_busy} cases {@link ResilientAgentLaunch} already retries inside one
* {@code agents.start} call — up to {@link #RELAUNCH_ATTEMPTS} times, sleeping via the
* injected sleeper between attempts, and returns the agent from the first attempt that
* succeeds.
*/
public Agent relaunch(String name) {
ResolvedLead resolved = resolveLaunchable(name);
if (resolved == null) {
return null;
}
for (int attempt = 1; attempt <= RELAUNCH_ATTEMPTS; attempt++) {
Agent started = launch(name, resolved.lead(), resolved.profile());
if (started != null) {
return started;
}
if (attempt < RELAUNCH_ATTEMPTS) {
sleeper.run();
}
}
return null;
}
/**
* How many live leads exist per configured name, and which of that name's labelled tabs are
* <em>not</em> live: a running agent in a tab labelled with that lead's exact {@code tab}
@@ -335,7 +404,7 @@ public final class LeadLauncher {
}
/**
* Start one lead. Returns false (having logged) rather than throwing on any failure.
* Start one lead. Returns null (having logged) rather than throwing on any failure.
*
* <p>Goes through the same {@link ResilientAgentLaunch} seam every member spawn uses
* (fleetd #727): the assembled argv is refused outright if it cannot fit the pane line herdr
@@ -344,7 +413,7 @@ public final class LeadLauncher {
* relaunch, and a seed pane whose shell has not reached its prompt yet ({@code
* agent_pane_busy}) is retried rather than failing on the first miss.
*/
private boolean launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
private Agent launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
String label = lead.tabLabel();
String cwd = (lead.cwd() == null || lead.cwd().isBlank())
? System.getProperty("user.dir") : lead.cwd();
@@ -375,7 +444,7 @@ public final class LeadLauncher {
log.info("lead '{}' launched: profile={} tab={} pane={} terminal={} label='{}' cwd={}",
name, profile.profile(), tab.tab().tabId(), started.paneId(),
started.terminalId(), label, cwd);
return true;
return started;
} catch (RuntimeException e) {
log.warn("lead '{}' failed to launch on profile '{}': {}",
name, profile.profile(), e.getMessage());
@@ -387,7 +456,7 @@ public final class LeadLauncher {
tab.tab().tabId(), cleanup.getMessage());
}
}
return false;
return null;
}
}
@@ -1,8 +1,11 @@
package dev.ltms.fleet.lead;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.WorkspaceControl;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -24,60 +27,69 @@ import java.util.function.Supplier;
/**
* fleetd #480: replace a lead session that has decided it is ready to be rolled over, without an
* operator doing it by hand. A lead writes a handover file, calls {@link #open}, and then — once
* every gate ({@link #confirm}'s own checks) has passed — a deferred, single-shot continuation
* clears the lead's own pane and bootstraps a fresh session against that file.
* every gate ({@link #confirm}'s own checks) has passed — a deferred, single-shot continuation ends
* the lead's own pane, launches a fresh one, and bootstraps that fresh session against the file.
*
* <p>This is the executor behind the {@code fleet_handover} MCP tool ({@code
* dev.ltms.fleet.mcp.FleetMcp#handover}), which drives {@link #open}, {@link #confirm}, {@link
* #cancel}, and {@link #status} from a tool call — wired in fleetd #480 Unit C. <strong>An earlier
* version of this paragraph said nothing called this class at all; that stopped being true once
* that unit landed, and this correction exists so the javadoc does not go on claiming it.</strong>
* #cancel}, and {@link #status} from a tool call.
*
* <p><strong>{@code confirm()} cannot roll inline — a fleetd #480 correction.</strong> The first
* version of this class called {@code agents.send(lead, "/clear")} directly from inside {@code
* confirm()}, then polled for the pane to become injectable again. That is wrong, because {@code
* confirm()} is called BY the lead, FROM the lead's own turn: the lead's pane is {@code WORKING}
* for the whole duration of that call and cannot possibly report injectable until {@code confirm()}
* itself returns. The poll always timed out — but only after the {@code /clear} had already been
* sent and queued in the pane, where it fired the instant the turn ended anyway. The result was the
* worst outcome this feature can produce: a silently destroyed lead context with no fresh session
* ever started, and a refusal return value that claimed nothing had happened.
*
* <p>The fix: {@link #confirm} validates every gate, then does no I/O against the lead's own pane
* at all — it only records that the request is approved and hands a one-shot continuation to
* {@code continuationRunner} before returning. That continuation is what actually touches the pane,
* once the calling turn has ended, in this order:
* <p><strong>{@code confirm()} cannot roll inline.</strong> {@code confirm()} is called BY the
* lead, FROM the lead's own turn: the lead's pane is {@code WORKING} for the whole duration of that
* call and cannot possibly report a real turn boundary until {@code confirm()} itself returns. So
* {@link #confirm} validates every gate, then does no I/O against the lead's own pane at all — it
* only records that the request is approved and hands a one-shot continuation to {@code
* continuationRunner} before returning. That continuation is what actually touches the pane, once
* the calling turn has ended, in this order:
* <ol>
* <li>wait for the lead's own pane to report a real turn boundary — {@code IDLE} or {@code
* DONE}, never merely {@code BLOCKED} — i.e. wait for the very {@code confirm()} call that
* approved this roll to finish its turn — bounded by {@code turnSettleSeconds}. <strong>If
* this never happens, nothing else in this list runs: no {@code /clear} is ever sent.</strong>
* A lead that never goes idle is a lead still doing real work, and clearing it would throw
* away live context — exactly the failure this correction exists to prevent.</li>
* <li>{@code agents.send(lead, "/clear")}</li>
* <li>wait for {@code /clear} to be picked up and settle, bounded by {@code clearSettleSeconds}
* (fleetd #489: no longer a plain re-check of the same boundary — {@code /clear} starts no
* turn of its own, so this instead nudges the submit keystroke while no pickup has been seen,
* then waits for a real {@code WORKING} → {@code IDLE}/{@code DONE} boundary once one has;
* see {@link #waitForClearPickupAndSettle})</li>
* <li>{@code agents.send(lead, cfg.bootstrapTextFor(p.handoverPath()))}</li>
* this never happens, nothing else in this list runs: the old pane is never touched.</strong>
* A lead that never goes idle is a lead still doing real work, and tearing it down would throw
* away live context.</li>
* <li>capture the old pane id (and, through it, the old tab) from {@link AgentControl#get}, with
* a bounded retry — the terminal-to-pane lookup it goes through can itself report a genuinely
* live agent as not found (see {@code AgentControl#agentCall}'s own re-resolve-once
* behaviour), and one false negative here must not abort an otherwise-healthy roll. Neither id
* is ever re-resolved from the terminal again after this — once the pane below is closed there
* is nothing left to resolve it from.</li>
* <li>resolve the lead's configured name from its terminal, for the relaunch step below.</li>
* <li>end the old session: close the pane (an already-gone pane counts as success; any other
* failure propagates), then close its tab only when the pane was that tab's sole occupant —
* the same pane-then-tab teardown {@code HerdrPeerLauncher#stop} uses for a member.</li>
* <li>confirm the old pane is actually gone by polling {@link
* dev.ltms.fleet.herdr.WorkspaceControl#locatePane} for a {@code null} result — never {@link
* AgentControl#status}, and never the live-lead terminal map, each of which answers a
* different question. <strong>If the old pane is never confirmed gone, no relaunch is
* attempted</strong> — see {@link RollState#OLD_PANE_NEVER_DIED}.</li>
* <li>launch a fresh lead with {@code LeadLauncher#relaunch}. <strong>If every attempt fails,
* {@code bootstrapText} is never sent</strong> — see {@link RollState#RELAUNCH_FAILED}.</li>
* <li>wait for the fresh pane to reach a real turn boundary ({@code IDLE} or {@code DONE},
* never merely {@code BLOCKED}), bounded by {@code relaunchReadySeconds}. This is the
* safety gate: typing into a pane that has not actually finished booting loses the
* keystrokes. <strong>If the pane never becomes ready, {@code bootstrapText} is never
* sent</strong> — see {@link RollState#RELAUNCH_NEVER_READY}.</li>
* <li>wait for the fresh terminal to be recognised as a live lead — present in the live-lead
* terminal map — bounded by {@code relaunchReadySeconds}. This is bookkeeping, not a
* safety gate: {@code bootstrapText} is sent either way once the pane is ready, whether or
* not this wait itself times out — see {@link RollState#RELAUNCH_NOT_RECOGNISED}.</li>
* <li>{@code agents.send(newTerminal, cfg.bootstrapTextFor(p.handoverPath()))} — sent to the
* FRESH terminal, never the one that was just torn down.</li>
* </ol>
* A {@link #confirm} that returns {@link RollDecision#approved()} therefore means <em>"every gate
* passed and the roll is scheduled"</em>, never <em>"the pane has been cleared"</em> — the pane may
* still be mid-turn, possibly for a long time, when the caller gets that answer back.
* passed and the roll is scheduled"</em>, never <em>"the lead has already been replaced"</em> — the
* old pane may still be mid-turn, possibly for a long time, when the caller gets that answer back.
*
* <p><strong>The safety invariant survives this change, restated precisely.</strong> The ticket
* that first defined this class required "no timer, no scheduler, no background thread" so that
* nothing but an explicit {@link #confirm} call could ever cause a {@code /clear}. That invariant
* is about INITIATIVE, not about synchronicity, and this correction keeps it: {@code
* continuationRunner} launches a single-shot task that exists only because one specific,
* <p><strong>The safety invariant.</strong> "No timer, no scheduler, no background thread" means
* that nothing but an explicit {@link #confirm} call can ever tear a lead's pane down.
* {@code continuationRunner} launches a single-shot task that exists only because one specific,
* already-approved {@link #confirm} call created it — it is not recurring, it is not started at
* construction time or on any schedule, and no two invocations of it ever share state. A recurring
* heartbeat or timer that could decide on its own initiative to roll a pane is still, and will
* always be, absent from this class. <strong>Nothing but an explicit {@link #confirm} call that
* passes every gate can ever cause a {@code /clear} — that call may simply finish its own work
* slightly later than the method return, as a continuation of the same approved request, rather
* than entirely inside the method body.</strong>
* heartbeat or timer that could decide on its own initiative to roll a pane is absent from this
* class. <strong>Nothing but an explicit {@link #confirm} call that passes every gate can ever tear
* a pane down — that call may simply finish its own work slightly later than the method return, as
* a continuation of the same approved request, rather than entirely inside the method body.</strong>
*
* <p><strong>Identity is resolved by the caller, never looked up here — a second fleetd #480
* correction.</strong> The first version resolved the pane to clear via {@code
@@ -115,19 +127,17 @@ public final class LeadRollover {
private static final Logger log = LoggerFactory.getLogger(LeadRollover.class);
/** Poll interval while waiting for the lead's pane to settle after {@code /clear}. */
static final long SETTLE_POLL_MS = 250;
/** Poll interval shared by every bounded wait in this class. */
static final long POLL_INTERVAL_MS = 250;
/**
* How many consecutive not-yet-picked-up polls {@link #waitForClearPickupAndSettle} allows
* before releasing rather than wedging the roll — the same constant and the same
* release-not-wedge choice {@link dev.ltms.fleet.inject.Injector} already makes for its own
* post-turn {@code /clear} housekeeping (fleetd #306). <strong>This bounds the number of
* consecutive polls, not the number of nudges:</strong> the first {@code PICKUP_GRACE_POLLS - 1}
* of those polls each send a nudge, and the {@code PICKUP_GRACE_POLLS}th releases instead of
* nudging again — so 8 polls produce 7 nudges, not 8.
* How long {@link #waitUntilPaneGone} polls {@link WorkspaceControl#locatePane} before giving
* up on ever seeing the old pane disappear. Not configurable: once {@link #endOldSession} has
* closed the pane (and, usually, its tab), herdr dropping the pane from its own bookkeeping is
* expected to show up within one or two polls, not on an operator-tunable timescale the way a
* CLI boot is.
*/
static final int PICKUP_GRACE_POLLS = 8;
static final int PANE_DEATH_TIMEOUT_SECONDS = 10;
/**
* One request opened by {@link #open}, pending its {@link #confirm} (or {@link #cancel}).
@@ -164,15 +174,21 @@ public final class LeadRollover {
* The handover file's modified time is not after {@link #open}'s request timestamp, or is
* older than {@code maxDocAgeSeconds}.
*/
HANDOVER_STALE
HANDOVER_STALE,
/**
* This lead terminal already has a roll running: an earlier {@link #confirm} call claimed
* it and that roll's continuation has not released it yet. {@code detail} names the lead
* terminal and the token that holds the claim.
*/
ROLL_ALREADY_RUNNING
}
/**
* The outcome of a {@link #confirm} call. {@link #approved()} means every gate passed and the
* roll has been handed to a one-shot continuation — <strong>not</strong> that the pane has been
* cleared; the continuation may still be waiting for the calling turn to end when this returns.
* Whether the deferred roll itself later goes on to clear the pane, refuse for never going
* idle, or refuse for never re-settling after {@code /clear} is logged only (see this class's
* roll has been handed to a one-shot continuation — <strong>not</strong> that the lead has
* already been replaced; the continuation may still be waiting for the calling turn to end when
* this returns. Whether the deferred roll itself later goes on to tear the old pane down and
* relaunch the lead, or refuses at any of its own steps, is logged only (see this class's
* javadoc) — there is deliberately no synchronous caller left by that point to hand a result to.
*/
public record RollDecision(boolean accepted, RefusalReason reason, String detail) {
@@ -208,7 +224,7 @@ public final class LeadRollover {
/**
* What is known about one token, right now — the answer {@link #status} gives. Distinguishes
* three terminal outcomes an approved roll can finish with, one in-flight outcome for a roll
* five terminal outcomes an approved roll can finish with, one in-flight outcome for a roll
* that has been approved but has not finished yet, and two answers for a token that names no
* active work at all: still pending confirmation, or nothing known about this token at all.
*/
@@ -231,41 +247,64 @@ public final class LeadRollover {
* #status} could wrongly answer {@link #UNKNOWN} ("nothing was ever requested") for a roll
* that is, in fact, actively running. This is not sticky: the deferred continuation
* overwrites this same entry with a terminal state ({@link #ROLLED}, {@link
* #TURN_NEVER_SETTLED}, {@link #CLEAR_NEVER_SETTLED}, or {@link #FAILED}) once it finishes
* — including by throwing, which fleetd #615's catch in {@link #runRollover} now turns into
* {@link #FAILED} instead of leaving this entry stuck forever.
* #TURN_NEVER_SETTLED}, {@link #OLD_PANE_NEVER_DIED}, {@link #RELAUNCH_FAILED}, {@link
* #RELAUNCH_NEVER_READY}, {@link #RELAUNCH_NOT_RECOGNISED}, or {@link #FAILED}) once it
* finishes — including by throwing, which {@link #runRollover}'s catch turns into {@link
* #FAILED} instead of leaving this entry stuck forever.
*/
IN_PROGRESS,
/**
* {@link #confirm} was approved and the deferred continuation completed the entire roll:
* the calling lead's turn settled, {@code /clear} was sent and settled, and {@code
* bootstrapText} was sent.
* {@link #confirm} was approved and the deferred continuation completed the entire roll: the
* calling lead's turn settled, the old pane was torn down and confirmed gone, a fresh lead
* was launched and recognised, and {@code bootstrapText} was sent to it.
*/
ROLLED,
/**
* {@link #confirm} was approved, but the calling lead's own turn never reached a boundary
* (IDLE or DONE) within {@code turnSettleSeconds} — no {@code /clear} was ever sent, at
* all. This is the branch the fleetd #480 correction exists to make safe, and the one this
* status exists to make VISIBLE: before this, a lead that hit this case had no way to find
* out, and would carry on believing it was about to be replaced. See this class's javadoc.
* (IDLE or DONE) within {@code turnSettleSeconds} — the old pane was never touched at all.
* This is the state that makes a lead's own stuck turn VISIBLE: without it, a lead that hit
* this case would have no way to find out, and would carry on believing it was about to be
* replaced. See this class's javadoc.
*/
TURN_NEVER_SETTLED,
/**
* {@link #confirm} was approved and {@code /clear} was sent, but the pane never re-settled
* within {@code clearSettleSeconds} — {@code bootstrapText} was never sent.
* {@link #confirm} was approved and the calling lead's turn settled, the old pane was closed
* (and its tab, if it was the sole occupant), but {@link
* dev.ltms.fleet.herdr.WorkspaceControl#locatePane} kept reporting it as still present for
* the whole pane-death timeout. No relaunch was ever attempted, and {@code bootstrapText}
* was never sent.
*/
CLEAR_NEVER_SETTLED,
OLD_PANE_NEVER_DIED,
/**
* fleetd #615: the deferred continuation threw a {@link RuntimeException} — most likely a
* {@link dev.ltms.fleet.herdr.HerdrException} out of one of the two unwrapped {@code
* agents.send} calls in {@link #runRollover} — and the continuation thread died with it.
* Before this state existed, that throw left {@link #outcomes} holding {@link #IN_PROGRESS}
* forever, because the production {@code continuationRunner} is a bare virtual thread with
* no uncaught-exception handler and nothing downstream of the throw ever ran to write a
* terminal outcome. {@code detail} names the exception, so a reader has something to act on
* — the same diagnostic style as {@link #TURN_NEVER_SETTLED} and {@link
* #CLEAR_NEVER_SETTLED}. The roll is dead at this point and does not retry itself; a stuck
* lead must {@link #open} a fresh request.
* The old pane was confirmed gone, but {@code LeadLauncher#relaunch} returned {@code null}
* — every launch attempt failed. {@code bootstrapText} was never sent, and no fresh terminal
* exists for this roll to have recognised.
*/
RELAUNCH_FAILED,
/**
* A fresh lead was launched, but its pane never reached a real turn boundary ({@code IDLE}
* or {@code DONE}, never merely {@code BLOCKED}) within {@code relaunchReadySeconds} — the
* CLI never finished booting, or it stayed paused on a startup prompt. {@code bootstrapText}
* was never sent: typing into a pane that is not actually ready to accept input loses the
* keystrokes.
*/
RELAUNCH_NEVER_READY,
/**
* A fresh lead was launched and its pane reached a real turn boundary, so {@code
* bootstrapText} WAS sent to it, but the terminal was never recognised as a live lead —
* present in the live-lead terminal map — within {@code relaunchReadySeconds}. The session
* itself is alive and bootstrapped; only the daemon's own bookkeeping has not caught up, and
* an operator should check why the tab was not recognised.
*/
RELAUNCH_NOT_RECOGNISED,
/**
* The deferred continuation threw a {@link RuntimeException} and the continuation thread
* died with it. Without this state, that throw would leave {@link #outcomes} holding {@link
* #IN_PROGRESS} forever, because the production {@code continuationRunner} is a bare virtual
* thread with no uncaught-exception handler and nothing downstream of the throw ever runs to
* write a terminal outcome. {@code detail} names the exception, so a reader has something to
* act on. The roll is dead at this point and does not retry itself; a stuck lead must
* {@link #open} a fresh request.
*/
FAILED,
/**
@@ -286,6 +325,10 @@ public final class LeadRollover {
public record RollStatus(RollState state, String detail) {}
private final AgentControl agents;
/** Workspace/tab/pane control — used to tear down the old pane and confirm it is gone. */
private final WorkspaceControl spaces;
/** Starts the fresh lead that replaces the one this roll tears down. */
private final LeadLauncher launcher;
private final Supplier<FleetConfig.LeadRollover> configSupplier;
/**
* Terminal id → that lead's configured workspace directory (their {@code
@@ -295,8 +338,21 @@ public final class LeadRollover {
* daemon-cwd bug this parameter exists to fix.
*/
private final Function<String, String> leadWorkspace;
/**
* Terminal id → that lead's configured name under {@code fleet.leaders}, or {@code null} when
* the terminal names no currently-recognised lead. The deferred continuation calls this, on the
* OLD terminal, before tearing it down, so it knows which lead to pass to {@link
* LeadLauncher#relaunch}.
*/
private final Function<String, String> leadNameForTerminal;
/**
* The daemon's current terminal id → lead name map, read fresh on every poll. The deferred
* continuation polls this for the FRESH terminal {@link LeadLauncher#relaunch} returns, to
* learn when that terminal has been recognised as a live lead — see this class's javadoc.
*/
private final Supplier<Map<String, String>> liveLeadTerminals;
private final LongSupplier nowMillis;
private final Runnable settleSleeper;
private final Runnable pollSleeper;
/**
* Launches the post-{@code confirm()} continuation. Production uses a single unstarted virtual
* thread per confirmed request — see this class's javadoc for why that is a single-shot task,
@@ -305,6 +361,15 @@ public final class LeadRollover {
*/
private final Consumer<Runnable> continuationRunner;
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
/**
* Lead terminal → the token of the roll currently holding that terminal exclusive, for
* {@link #confirm}'s single-flight claim. {@link #confirm} claims an entry here with an
* atomic put-if-absent once every other gate has passed, refusing with {@link
* RefusalReason#ROLL_ALREADY_RUNNING} when a claim is already held; {@link #runRollover}
* releases it in a {@code finally}, on both the success and the thrown-exception path. A
* terminal absent from this map has no roll currently in flight for it.
*/
private final Map<String, String> rollingByTerminal = new ConcurrentHashMap<>();
/**
* Finished tokens → what actually happened, for {@link #status}. Bounded by {@link
* #OUTCOME_HISTORY_CAP}, oldest evicted first ({@code removeEldestEntry} on an insertion-order
@@ -323,29 +388,40 @@ public final class LeadRollover {
}
});
/** Production constructor — wall clock, real sleep between settle polls, a real virtual thread. */
public LeadRollover(AgentControl agents, Supplier<FleetConfig.LeadRollover> configSupplier,
Function<String, String> leadWorkspace) {
this(agents, configSupplier, leadWorkspace, System::currentTimeMillis,
() -> sleepUninterruptibly(SETTLE_POLL_MS),
/** Production constructor — wall clock, real sleep between polls, a real virtual thread. */
public LeadRollover(AgentControl agents, WorkspaceControl spaces, LeadLauncher launcher,
Supplier<FleetConfig.LeadRollover> configSupplier,
Function<String, String> leadWorkspace,
Function<String, String> leadNameForTerminal,
Supplier<Map<String, String>> liveLeadTerminals) {
this(agents, spaces, launcher, configSupplier, leadWorkspace, leadNameForTerminal,
liveLeadTerminals, System::currentTimeMillis,
() -> sleepUninterruptibly(POLL_INTERVAL_MS),
r -> Thread.ofVirtual().name("lead-rollover-continuation-").start(r));
}
/**
* Full constructor — an injectable wall-clock supplier, settle-poll sleeper, and continuation
* runner, for tests. {@code nowMillis} MUST be a wall-clock source (e.g. {@code
* Full constructor — an injectable wall-clock supplier, poll sleeper, and continuation runner,
* for tests. {@code nowMillis} MUST be a wall-clock source (e.g. {@code
* System.currentTimeMillis()}), never {@code System.nanoTime()}: the freshness check compares
* against a file's modified time, which only a wall clock is comparable to, and {@code
* nanoTime} freezes while the host sleeps (fleetd #386).
* nanoTime} freezes while the host sleeps.
*/
LeadRollover(AgentControl agents, Supplier<FleetConfig.LeadRollover> configSupplier,
Function<String, String> leadWorkspace, LongSupplier nowMillis,
Runnable settleSleeper, Consumer<Runnable> continuationRunner) {
LeadRollover(AgentControl agents, WorkspaceControl spaces, LeadLauncher launcher,
Supplier<FleetConfig.LeadRollover> configSupplier,
Function<String, String> leadWorkspace,
Function<String, String> leadNameForTerminal,
Supplier<Map<String, String>> liveLeadTerminals,
LongSupplier nowMillis, Runnable pollSleeper, Consumer<Runnable> continuationRunner) {
this.agents = agents;
this.spaces = spaces;
this.launcher = launcher;
this.configSupplier = configSupplier;
this.leadWorkspace = leadWorkspace;
this.leadNameForTerminal = leadNameForTerminal;
this.liveLeadTerminals = liveLeadTerminals;
this.nowMillis = nowMillis;
this.settleSleeper = settleSleeper;
this.pollSleeper = pollSleeper;
this.continuationRunner = continuationRunner;
}
@@ -487,6 +563,16 @@ public final class LeadRollover {
return docCheck;
}
// Single-flight claim: atomic put-if-absent, taken only after every other gate has
// passed, so a refused confirm() never takes it. A non-null previous value means a
// different, still-running roll already holds this lead terminal.
String holder = rollingByTerminal.putIfAbsent(p.leadTerminal(), token);
if (holder != null) {
return RollDecision.refused(RefusalReason.ROLL_ALREADY_RUNNING,
"lead terminal " + p.leadTerminal() + " already has a roll running under token "
+ holder);
}
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
// OUTCOME_HISTORY_CAP's javadoc. This ordering means `token` is written into `outcomes`
// while it is STILL present in `pending`; status() checks `outcomes` first (see that
@@ -494,34 +580,48 @@ public final class LeadRollover {
// remove-then-put ordering would leave in which the token is in neither map.
outcomes.put(token, new RollStatus(RollState.IN_PROGRESS,
"confirm() approved this roll and handed it to the deferred continuation; it has "
+ "not finished yet — still waiting for the calling turn to settle, for "
+ "/clear to be sent and settle, or for bootstrapText to be sent"));
+ "not finished yet — still waiting for the calling turn to settle, for the "
+ "old pane to be torn down and confirmed gone, for the fresh lead to be "
+ "recognised, or for bootstrapText to be sent"));
pending.remove(token);
log.info("lead-rollover: confirmed token={} lead={} — roll scheduled once the calling turn ends",
token, callerTerminal);
continuationRunner.accept(() -> runRollover(p, cfg));
try {
continuationRunner.accept(() -> runRollover(p, cfg));
} catch (RuntimeException e) {
// continuationRunner can reject the hand-off itself (e.g. a bounded executor's
// RejectedExecutionException) before runRollover ever starts, so runRollover's own
// finally — the only other place that releases rollingByTerminal — never runs either.
// Release the claim here and overwrite the IN_PROGRESS entry with a terminal outcome,
// or this lead terminal could never be rolled again and status() would report
// IN_PROGRESS forever for a roll that in fact never started.
log.warn("lead-rollover: continuationRunner rejected token={} lead={}: {} — the roll "
+ "never started; releasing its claim and reporting it as FAILED",
token, callerTerminal, e.toString(), e);
rollingByTerminal.remove(p.leadTerminal(), token);
outcomes.put(token, new RollStatus(RollState.FAILED,
"continuationRunner rejected this roll before it ever started: " + e.toString()
+ " — the roll never ran; open() a fresh rollover request"));
}
return RollDecision.approved();
}
/**
* The single-shot continuation {@link #confirm} hands to {@code continuationRunner}. Runs
* entirely after {@link #confirm} has returned to its caller — see this class's javadoc for the
* four-step order. There is no result to return to by this point, so every outcome is logged
* only.
* full order. There is no result to return to by this point, so every outcome is logged only.
*
* <p><strong>fleetd #615 — the whole body is wrapped in one {@code try}.</strong> The two {@code
* agents.send} calls below are not wrapped individually: {@code send} → {@code agentCall} →
* {@code herdr.call} can throw an unchecked {@link dev.ltms.fleet.herdr.HerdrException} (see
* {@code AgentControl.java}), and the production {@code continuationRunner} is a bare virtual
* thread with no uncaught-exception handler (see this class's public constructor). Before this
* fix, either throw killed the continuation thread silently, leaving the {@link
* RollState#IN_PROGRESS} entry {@link #confirm} wrote at hand-off stuck forever — {@link
* #status} had no way to tell a dead roll from one still genuinely running. The {@code catch}
* below is scoped to the method body rather than to each {@code send} call individually, so it
* also covers anything else added to this continuation later, not just today's two call sites —
* the same reasoning that put the write-a-terminal-outcome step at each of this method's other
* exits (see the {@link RollState#TURN_NEVER_SETTLED} and {@link RollState#CLEAR_NEVER_SETTLED}
* branches below) rather than inside the helpers that detect them.</p>
* <p><strong>The whole body is wrapped in one {@code try}.</strong> Several calls below —
* {@code agents.get}, {@code agents.close}, {@code agents.send} — can throw an unchecked {@link
* dev.ltms.fleet.herdr.HerdrException} (see {@code AgentControl.java}), and the production
* {@code continuationRunner} is a bare virtual thread with no uncaught-exception handler (see
* this class's public constructor). An uncaught throw would kill the continuation thread
* silently, leaving the {@link RollState#IN_PROGRESS} entry {@link #confirm} wrote at hand-off
* stuck forever — {@link #status} would have no way to tell a dead roll from one still
* genuinely running. The {@code catch} below is scoped to the method body rather than to each
* call individually, so it also covers every call in this continuation, not a fixed list of
* call sites — the same reasoning that put the write-a-terminal-outcome step at each of this
* method's other exits rather than inside the helpers that detect them.</p>
*
* <p>Only {@link RuntimeException} is caught, matching the local convention {@link
* #waitUntilAtTurnBoundary} already set around its own {@code agents.status} call — not the
@@ -539,6 +639,12 @@ public final class LeadRollover {
"the roll's continuation threw " + e.toString() + " — the roll is dead and will "
+ "not retry itself; check the daemon log for the stack trace, then open() "
+ "a fresh rollover request"));
} finally {
// Release the single-flight claim on both the normal return and the thrown-exception
// path above — a release only on success would leave this lead terminal unrollable
// forever after one failure. The conditional two-argument remove only clears the
// entry this roll itself holds, never a different roll's claim on the same terminal.
rollingByTerminal.remove(p.leadTerminal(), p.token());
}
}
@@ -546,53 +652,237 @@ public final class LeadRollover {
private void runRolloverUnguarded(PendingRollover p, FleetConfig.LeadRollover cfg) {
String lead = p.leadTerminal();
long rollStartMillis = nowMillis.getAsLong();
TurnSettleResult turnResult = waitUntilAtTurnBoundary(lead, cfg.turnSettleSeconds());
if (!turnResult.settled()) {
// fleetd #494 follow-up: this line had the SAME defect as the /clear-timeout line below
// — cfg.turnSettleSeconds() is the CONFIGURED budget, not how long this wait actually
// ran. Print the measured elapsed time alongside it, labelled, exactly like the /clear
// path already does.
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) after "
+ "confirm() — refusing to send /clear at all; the calling lead's own "
+ "turn is still live and clearing it now would destroy live context "
+ "confirm() — the old pane is never touched; the calling lead's own "
+ "turn is still live and tearing it down now would destroy live context "
+ "(token={}, configured={}s elapsed={}ms)",
lead, p.token(), cfg.turnSettleSeconds(), turnResult.elapsedMillis());
outcomes.put(p.token(), new RollStatus(RollState.TURN_NEVER_SETTLED,
"the calling lead's own turn never reached a boundary (IDLE or DONE) within "
+ "turnSettleSeconds=" + cfg.turnSettleSeconds() + "s (measured elapsed="
+ turnResult.elapsedMillis() + "ms) — no /clear was ever sent. If this "
+ "keeps happening, raise turnSettleSeconds in fleetd.yaml"));
+ turnResult.elapsedMillis() + "ms) — the old pane was never touched. If "
+ "this keeps happening, raise turnSettleSeconds in fleetd.yaml"));
return;
}
// This deliberately bypasses Injector, exactly like ClaudeCodeLauncher#clearContext:
// /clear is housekeeping, not a delegated turn, and routing it through Injector wedges the
// pane forever (see this class's javadoc).
agents.send(lead, "/clear");
ClearSettleResult clearResult = waitForClearPickupAndSettle(lead, cfg.clearSettleSeconds());
if (!clearResult.settled()) {
// fleetd #494: cfg.clearSettleSeconds() is the CONFIGURED budget, not how long the wait
// actually ran — an operator reading only that number wrongly believes it is a measured
// duration. Print the measured elapsed time and nudge count alongside it, each labelled,
// so the two can be compared at a glance.
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) after "
+ "/clear — NOT sending bootstrapText (token={}, configured={}s "
+ "elapsed={}ms nudges={})",
lead, p.token(), cfg.clearSettleSeconds(), clearResult.elapsedMillis(),
clearResult.nudges());
outcomes.put(p.token(), new RollStatus(RollState.CLEAR_NEVER_SETTLED,
"/clear was sent, but the pane never re-settled within clearSettleSeconds="
+ cfg.clearSettleSeconds() + "s (measured elapsed=" + clearResult.elapsedMillis()
+ "ms, nudges=" + clearResult.nudges() + ") — bootstrapText was never sent"));
// Captured once, here, and never re-resolved from `lead` again below: once the pane is
// closed there is nothing left for a terminal lookup to find.
Agent oldAgent = captureAgentWithRetry(lead);
String oldPaneId = oldAgent.paneId();
String leadName = leadNameForTerminal.apply(lead);
endOldSession(oldPaneId);
DeathResult deathResult = waitUntilPaneGone(oldPaneId);
if (!deathResult.gone()) {
log.warn("lead-rollover: old pane {} for lead {} was never confirmed gone after being "
+ "closed — not attempting a relaunch (token={}, timeout={}s "
+ "elapsed={}ms)",
oldPaneId, lead, p.token(), PANE_DEATH_TIMEOUT_SECONDS, deathResult.elapsedMillis());
outcomes.put(p.token(), new RollStatus(RollState.OLD_PANE_NEVER_DIED,
"the old pane was closed, but locatePane kept reporting it as still present "
+ "after a pane-death timeout=" + PANE_DEATH_TIMEOUT_SECONDS
+ "s (measured elapsed=" + deathResult.elapsedMillis() + "ms) — no "
+ "relaunch was attempted"));
return;
}
agents.send(lead, cfg.bootstrapTextFor(p.handoverPath()));
Agent newAgent = launcher.relaunch(leadName);
if (newAgent == null) {
log.warn("lead-rollover: relaunch of lead '{}' (old terminal {}) failed every attempt "
+ "— bootstrapText was never sent (token={})", leadName, lead, p.token());
outcomes.put(p.token(), new RollStatus(RollState.RELAUNCH_FAILED,
"lead '" + leadName + "' could not be relaunched — every attempt failed; "
+ "bootstrapText was never sent"));
return;
}
ReadinessResult readinessResult = waitUntilPaneReady(newAgent.terminalId(),
cfg.relaunchReadySeconds());
if (!readinessResult.ready()) {
log.warn("lead-rollover: fresh pane for lead '{}' (terminal {}) never reached a real "
+ "turn boundary — bootstrapText was never sent (token={}, configured={}s "
+ "elapsed={}ms)",
leadName, newAgent.terminalId(), p.token(), cfg.relaunchReadySeconds(),
readinessResult.elapsedMillis());
outcomes.put(p.token(), new RollStatus(RollState.RELAUNCH_NEVER_READY,
"fresh terminal " + newAgent.terminalId() + " never reached a real turn "
+ "boundary (IDLE or DONE) within relaunchReadySeconds="
+ cfg.relaunchReadySeconds() + "s (measured elapsed="
+ readinessResult.elapsedMillis() + "ms) — bootstrapText was never "
+ "sent"));
return;
}
IdentityResult identityResult = waitUntilRecognisedAsLead(newAgent.terminalId(),
cfg.relaunchReadySeconds());
agents.send(newAgent.terminalId(), cfg.bootstrapTextFor(p.handoverPath()));
if (!identityResult.ready()) {
log.warn("lead-rollover: fresh terminal {} for lead '{}' is alive and bootstrapped, but "
+ "was never recognised as a live lead — an operator should check why "
+ "the tab was not recognised (token={}, configured={}s elapsed={}ms)",
newAgent.terminalId(), leadName, p.token(), cfg.relaunchReadySeconds(),
identityResult.elapsedMillis());
outcomes.put(p.token(), new RollStatus(RollState.RELAUNCH_NOT_RECOGNISED,
"bootstrapText was sent to fresh terminal " + newAgent.terminalId() + ", but "
+ "it was never recognised as a live lead within relaunchReadySeconds="
+ cfg.relaunchReadySeconds() + "s (measured elapsed="
+ identityResult.elapsedMillis() + "ms) — check why the tab was not "
+ "recognised"));
return;
}
long rollElapsedMillis = nowMillis.getAsLong() - rollStartMillis;
log.info("lead-rollover: rolled token={} lead={} elapsedMs={}", p.token(), lead, rollElapsedMillis);
log.info("lead-rollover: rolled token={} oldLead={} newTerminal={} elapsedMs={}",
p.token(), lead, newAgent.terminalId(), rollElapsedMillis);
outcomes.put(p.token(), new RollStatus(RollState.ROLLED,
"rolled successfully in " + rollElapsedMillis + "ms"));
"rolled successfully in " + rollElapsedMillis + "ms; new terminal="
+ newAgent.terminalId()));
}
/** Attempts {@link #captureAgentWithRetry} makes before letting the failure propagate. */
static final int CAPTURE_RETRIES = 3;
/**
* {@link AgentControl#get} for {@code lead}, retried up to {@link #CAPTURE_RETRIES} times. The
* terminal-to-pane lookup it goes through can report a genuinely live agent as not found (see
* {@code AgentControl#agentCall}'s own re-resolve-once behaviour), and one such false negative
* must not abort an otherwise-healthy roll. The result is captured once by the caller and never
* looked up again — see this class's javadoc.
*
* @throws RuntimeException the last failure, if every attempt fails — {@link #runRollover}'s
* catch turns that into {@link RollState#FAILED}
*/
private Agent captureAgentWithRetry(String lead) {
RuntimeException last = null;
for (int attempt = 1; attempt <= CAPTURE_RETRIES; attempt++) {
try {
return agents.get(lead);
} catch (RuntimeException e) {
last = e;
log.debug("lead-rollover: agents.get({}) failed on attempt {}/{}: {}",
lead, attempt, CAPTURE_RETRIES, e.toString());
if (attempt < CAPTURE_RETRIES) {
pollSleeper.run();
}
}
}
throw last;
}
/**
* End the old lead's session: close its pane, then close its tab only when the pane was that
* tab's sole occupant — the same pane-then-tab teardown {@code HerdrPeerLauncher#stop} uses for
* a member. An already-gone pane counts as success; any other {@code agents.close} failure
* propagates, so a genuinely failed teardown is never reported as done. A failing
* {@code spaces.closeTab} never propagates — by the time it runs the pane is already closed, so
* it is cosmetic tidying, not a real teardown failure.
*/
private void endOldSession(String paneId) {
WorkspaceControl.PaneLocation loc = spaces.locatePane(paneId);
try {
agents.close(paneId);
} catch (HerdrException e) {
if (!isAlreadyGone(e)) {
throw e;
}
log.debug("lead-rollover: pane.close({}) ignored — already gone: {}", paneId, e.getMessage());
}
if (loc != null && loc.tabPaneCount() == 1) {
try {
spaces.closeTab(loc.tabId());
} catch (RuntimeException e) {
log.warn("lead-rollover: tab.close({}) failed — the pane is already torn down, so "
+ "continuing; the tab may need manual cleanup: {}", loc.tabId(), e.getMessage());
}
} else if (loc != null) {
log.debug("lead-rollover: not closing tab {} — it holds {} panes (not a dedicated lead "
+ "tab)", loc.tabId(), loc.tabPaneCount());
}
}
/** True when a herdr error means the target is already gone (safe to treat as done). */
private static boolean isAlreadyGone(HerdrException e) {
return e.code() != null && e.code().endsWith("_not_found");
}
/**
* Poll {@link WorkspaceControl#locatePane} for {@code paneId} until it reports {@code null}
* (the pane is gone) or {@link #PANE_DEATH_TIMEOUT_SECONDS} elapses. Deliberately never calls
* {@link AgentControl#status} and never reads the live-lead terminal map — both answer a
* different question (whether an AGENT is live, not whether this PANE still exists) and
* {@code locatePane} alone catches a {@link HerdrException} from the underlying {@code
* pane.get} and turns it into {@code null} — see this class's javadoc.
*/
private DeathResult waitUntilPaneGone(String paneId) {
long startMillis = nowMillis.getAsLong();
long deadline = startMillis + TimeUnit.SECONDS.toMillis(PANE_DEATH_TIMEOUT_SECONDS);
while (nowMillis.getAsLong() < deadline) {
if (spaces.locatePane(paneId) == null) {
return new DeathResult(true, nowMillis.getAsLong() - startMillis);
}
pollSleeper.run();
}
return new DeathResult(false, nowMillis.getAsLong() - startMillis);
}
/** The measured outcome of {@link #waitUntilPaneGone}. */
private record DeathResult(boolean gone, long elapsedMillis) {}
/**
* Poll until {@code newTerminal}'s own pane reaches a real turn boundary ({@link
* AgentStatus#IDLE} or {@link AgentStatus#DONE}, never merely {@link AgentStatus#BLOCKED}) —
* the same exclusion {@link #waitUntilAtTurnBoundary} applies to the calling lead's own turn,
* applied here to the fresh one, so {@code bootstrapText} is never typed into a pane that has
* not actually finished booting — or {@code readySeconds} elapses. A failed status read
* degrades to "not yet ready" and is retried on the next poll.
*/
private ReadinessResult waitUntilPaneReady(String newTerminal, int readySeconds) {
long startMillis = nowMillis.getAsLong();
long deadline = startMillis + TimeUnit.SECONDS.toMillis(readySeconds);
while (nowMillis.getAsLong() < deadline) {
AgentStatus status;
try {
status = agents.status(newTerminal);
} catch (RuntimeException e) {
log.debug("lead-rollover: status check failed while waiting for {} to be ready: {}",
newTerminal, e.toString());
status = null;
}
if (status == AgentStatus.IDLE || status == AgentStatus.DONE) {
return new ReadinessResult(true, nowMillis.getAsLong() - startMillis);
}
pollSleeper.run();
}
return new ReadinessResult(false, nowMillis.getAsLong() - startMillis);
}
/** The measured outcome of {@link #waitUntilPaneReady}. */
private record ReadinessResult(boolean ready, long elapsedMillis) {}
/**
* Poll until {@code newTerminal} is present in {@link #liveLeadTerminals} or {@code
* readySeconds} elapses. This is bookkeeping, not a safety gate: the pane's own readiness (see
* {@link #waitUntilPaneReady}) is what decides whether {@code bootstrapText} is safe to send —
* a timeout here only means the daemon's own lead-discovery scan has not caught up yet.
*/
private IdentityResult waitUntilRecognisedAsLead(String newTerminal, int readySeconds) {
long startMillis = nowMillis.getAsLong();
long deadline = startMillis + TimeUnit.SECONDS.toMillis(readySeconds);
while (nowMillis.getAsLong() < deadline) {
if (liveLeadTerminals.get().containsKey(newTerminal)) {
return new IdentityResult(true, nowMillis.getAsLong() - startMillis);
}
pollSleeper.run();
}
return new IdentityResult(false, nowMillis.getAsLong() - startMillis);
}
/** The measured outcome of {@link #waitUntilRecognisedAsLead}. */
private record IdentityResult(boolean ready, long elapsedMillis) {}
/** Drop a pending request without rolling. @return whether a pending request existed for {@code token} */
public boolean cancel(String token) {
return pending.remove(token) != null;
@@ -682,15 +972,12 @@ public final class LeadRollover {
/**
* Poll {@link AgentControl#status} until {@code target} reports a real turn boundary — {@link
* AgentStatus#IDLE} or {@link AgentStatus#DONE} — bounded by {@code settleSeconds}. Used once by
* {@link #runRollover}, to wait for the CALLING turn's own pane to settle before {@code /clear}
* is ever sent at all — the {@code turnSettleSeconds} gate that makes this correction safe. The
* SECOND wait, after {@code /clear}, is {@link #waitForClearPickupAndSettle} instead (fleetd
* #489) — a plain boundary check is not enough there, because {@code /clear} starts no turn of
* its own, so this method would (wrongly) report "settled" on its very first poll whether or not
* {@code /clear} was actually picked up. A failed status read degrades to "not yet settled" and
* is retried on the next poll, the same posture {@code LeadHeartbeatLoop} and {@code
* HerdrPeerLauncher}'s readiness gate already take toward an unreadable status.
* AgentStatus#IDLE} or {@link AgentStatus#DONE} — bounded by {@code settleSeconds}. Used by
* {@link #runRollover} to wait for the CALLING turn's own pane to settle before the old pane is
* touched at all — the {@code turnSettleSeconds} gate that makes tearing it down safe. A failed
* status read degrades to "not yet settled" and is retried on the next poll, the same posture
* {@code LeadHeartbeatLoop} and {@code HerdrPeerLauncher}'s readiness gate already take toward
* an unreadable status.
*
* <p><strong>Deliberately not {@link AgentStatus#injectable()}.</strong> {@code injectable()}
* answers the {@code Injector}'s question — "may I deliver a message without stepping on a live
@@ -698,16 +985,16 @@ public final class LeadRollover {
* an approval prompt is safe to queue a message behind. This class asks a stricter question —
* "has the turn actually ended" — and {@code BLOCKED} answers no: it is a live turn that is
* merely paused, not one that has finished. Reusing {@code injectable()} here would let this
* wait fire {@code /clear} while the lead's own {@code confirm()}-calling turn is still live and
* paused on a prompt — exactly the live-context-destroying failure the {@code turnSettleSeconds}
* gate exists to prevent. Do not "simplify" this back to {@code injectable()}. ({@link
* #waitForClearPickupAndSettle} keeps the same exclusion of {@code BLOCKED}, for the same
* reason, on the second wait.)
* wait tear the old pane down while the lead's own {@code confirm()}-calling turn is still live
* and paused on a prompt — exactly the live-context-destroying failure {@code turnSettleSeconds}
* exists to prevent. Do not "simplify" this back to {@code injectable()}. ({@link
* #waitUntilPaneReady} applies the same exclusion of {@code BLOCKED} to the fresh lead's own
* turn.)
*
* @return a {@link TurnSettleResult} whose {@code settled()} is {@code true} once a real
* boundary was observed, {@code false} if {@code settleSeconds} elapses first.
* {@code elapsedMillis()} is a MEASURED value from the injected {@link #nowMillis}
* clock, never the configured {@code settleSeconds} budget (fleetd #494 follow-up).
* clock, never the configured {@code settleSeconds} budget.
*/
private TurnSettleResult waitUntilAtTurnBoundary(String target, int settleSeconds) {
long startMillis = nowMillis.getAsLong();
@@ -724,142 +1011,11 @@ public final class LeadRollover {
if (status == AgentStatus.IDLE || status == AgentStatus.DONE) {
return new TurnSettleResult(true, nowMillis.getAsLong() - startMillis);
}
settleSleeper.run();
pollSleeper.run();
}
return new TurnSettleResult(false, nowMillis.getAsLong() - startMillis);
}
/**
* The measured outcome of {@link #waitUntilAtTurnBoundary} — fleetd #494 follow-up. The sibling
* of {@link ClearSettleResult} for the FIRST wait, which never nudges, so it carries no nudge
* count.
*/
/** The measured outcome of {@link #waitUntilAtTurnBoundary}. */
private record TurnSettleResult(boolean settled, long elapsedMillis) {}
/**
* The SECOND wait in {@link #runRollover} — after {@code /clear} has been sent, waits for it to
* settle, bounded by {@code settleSeconds}. <strong>fleetd #489 — the paste-race fix.</strong>
* {@code /clear} does not start a real turn of its own, so a pane with no submit race simply
* stays {@link AgentStatus#IDLE} the whole time: {@link #waitUntilAtTurnBoundary} would (wrongly)
* call that "settled" on its very first poll, whether or not the {@code /clear} Enter actually
* landed. That was Fault 1, measured live on 2026-09-12 — the second gate was a no-op, so a
* {@code bootstrapText} send followed immediately, racing Fault 2: {@link AgentControl#submit}'s
* own javadoc already records that the submit accompanying a delivery "can race the paste —
* especially right as the worker's TUI becomes interactive — leaving the text unsubmitted"
* (CB-113). Because {@code runRollover} deliberately bypasses {@code Injector} for {@code
* /clear} (see this class's javadoc), it inherited none of {@code Injector}'s nudging — so the
* lost {@code /clear} Enter sat in the input box and {@code bootstrapText} was typed right after
* it, landing as one concatenated line.
*
* <p>This method copies the pickup-nudge pattern {@link dev.ltms.fleet.inject.Injector} already
* ships for exactly this, on its own post-turn {@code /clear} housekeeping (fleetd #306; see
* {@code Injector.java:288-340} and {@code Injector.java:437-442}):
* <ul>
* <li>an {@link AgentStatus#WORKING} sample means {@code /clear} was picked up as a real
* turn;</li>
* <li>until that happens, each poll that still reports {@link AgentStatus#IDLE} or {@link
* AgentStatus#DONE} re-sends the submit keystroke ({@link AgentControl#submit}) to nudge
* the raced Enter — for the first {@code PICKUP_GRACE_POLLS - 1} of {@link
* #PICKUP_GRACE_POLLS} consecutive such polls (i.e. {@code PICKUP_GRACE_POLLS - 1}
* nudges: 7, not 8, given {@code PICKUP_GRACE_POLLS = 8}). A second Enter on an empty
* Claude Code prompt is a no-op, so repeating it is safe;</li>
* <li>the {@code PICKUP_GRACE_POLLS}th consecutive such poll, with {@code WORKING} still never
* observed, releases rather than wedges the roll instead of nudging again — the same
* choice {@code Injector} makes — and returns {@code settled() == true} anyway, logged at
* {@code warn} with the measured elapsed time (fleetd #494) so an operator can see which
* path ran and how long it actually took;</li>
* <li>once {@code WORKING} has been observed, nudging stops and this instead waits for a real
* {@code working → IDLE/DONE} completion boundary before returning {@code true}.</li>
* </ul>
*
* <p><strong>{@link AgentStatus#BLOCKED} is deliberately excluded from both the nudge and the
* boundary check</strong> — the same reasoning as {@link #waitUntilAtTurnBoundary}'s own
* javadoc: a paused live turn is not a settled one, and re-sending Enter into an open approval
* prompt could wrongly answer it. A {@code BLOCKED} sample (or an unreadable/{@link
* AgentStatus#UNKNOWN} one) simply keeps this polling, with no nudge and no release, until either
* a real boundary is reached or {@code settleSeconds} runs out.
*
* <p>{@link AgentControl#submit} can itself throw; a {@link RuntimeException} from it is
* swallowed and logged at {@code debug}, exactly like {@code Injector.java:437-442} — a failed
* nudge must not abort the roll.
*
* @return a {@link ClearSettleResult} whose {@code settled()} is {@code true} once {@code
* /clear} has settled, or once the nudge budget was exhausted with no pickup ever
* observed (released rather than wedged); {@code false} if {@code settleSeconds} elapses
* first — the caller must NOT send {@code bootstrapText} in that case, exactly as before
* this fix. {@code elapsedMillis()} and {@code nudges()} are MEASURED values (from the
* injected {@link #nowMillis} clock and an actual nudge count), never the configured
* {@code settleSeconds} budget (fleetd #494).
*/
private ClearSettleResult waitForClearPickupAndSettle(String target, int settleSeconds) {
long startMillis = nowMillis.getAsLong();
long deadline = startMillis + TimeUnit.SECONDS.toMillis(settleSeconds);
boolean pickedUp = false; // a WORKING sample has been observed since /clear was sent
int idlePollsAwaitingPickup = 0;
int nudges = 0;
while (nowMillis.getAsLong() < deadline) {
AgentStatus status;
try {
status = agents.status(target);
} catch (RuntimeException e) {
log.debug("lead-rollover: status check failed while waiting for {} to settle after "
+ "/clear: {}", target, e.toString());
status = null;
}
if (status == AgentStatus.WORKING) {
pickedUp = true;
} else if (status == AgentStatus.IDLE || status == AgentStatus.DONE) {
if (pickedUp) {
// a real WORKING -> IDLE/DONE completion boundary
return new ClearSettleResult(true, nowMillis.getAsLong() - startMillis, nudges);
}
if (++idlePollsAwaitingPickup >= PICKUP_GRACE_POLLS) {
long elapsedMillis = nowMillis.getAsLong() - startMillis;
// fleetd #494: this release trades a possibly-unsubmitted /clear for progress
// instead of wedging the roll — that trade is deliberate and stays. But it is
// also exactly the case that reported false success in the real incident (the
// whole roll "succeeded" after 438ms of a 20s budget), so raise it to WARN and
// print the MEASURED elapsed time next to the target pane, not just the count.
//
// fleetd #494 follow-up (2nd pass): BOTH numbers in this line must come from
// the loop's own counters, never from the PICKUP_GRACE_POLLS constant.
// `idlePollsAwaitingPickup` and `nudges` each have exactly one write site in
// this loop, on the same branch, so on this branch they cannot differ from
// PICKUP_GRACE_POLLS / PICKUP_GRACE_POLLS - 1 today — no test can prove the
// difference on this line, and printing the counters does not change that.
// What it does buy: one source of truth instead of two, so a later change to
// the loop (an early return, a second increment site, a different exit
// condition) cannot leave this message reporting a number the loop no longer
// produces. The place where `nudges` genuinely varies with the run — and is
// covered by a test that can tell it apart from a constant — is the
// /clear-timeout warn in runRollover, which prints clearResult.nudges().
log.warn("lead-rollover: /clear on {} was never observed as WORKING after {} "
+ "consecutive IDLE/DONE polls ({} of those were nudged) — "
+ "releasing rather than wedging the roll (elapsed={}ms)",
target, idlePollsAwaitingPickup, nudges, elapsedMillis);
return new ClearSettleResult(true, elapsedMillis, nudges);
}
try {
agents.submit(target); // nudge a raced Enter (CB-113) so /clear actually submits
} catch (RuntimeException e) {
log.debug("lead-rollover: resubmit to {} failed (will retry next poll): {}",
target, e.getMessage());
} finally {
nudges++; // an attempted nudge, whether or not the submit call itself threw
}
}
// AgentStatus.BLOCKED or UNKNOWN (or an unreadable status, above): neither a pickup
// signal nor a boundary — keep polling without nudging or releasing.
settleSleeper.run();
}
return new ClearSettleResult(false, nowMillis.getAsLong() - startMillis, nudges);
}
/**
* The measured outcome of {@link #waitForClearPickupAndSettle} — fleetd #494. Carries the
* MEASURED elapsed time (from the injected {@link #nowMillis} clock) and nudge count alongside
* the settle/timeout decision, so callers can log them instead of the configured budget, which
* is not how long the wait actually ran.
*/
private record ClearSettleResult(boolean settled, long elapsedMillis, int nudges) {}
}
@@ -334,7 +334,7 @@ public final class FleetMcp {
* caller explicitly saying so — never by omitting a {@link CallerResolver} the way the old
* {@code callers == null} idiom allowed. {@code callers} itself is required either way: even
* under {@link #UNENFORCED}, the one real {@link CallerResolver} still resolves every caller's
* {@link Principal} (so {@code markSpawnedMemberPresent}/{@code recordPrimarySingleton} see a
* {@link Principal} (so {@code markTrackedCallerPresent}/{@code recordPrimarySingleton} see a
* real identity), and {@link #denyFor} is the only thing that changes.
*/
public enum AuthorizationMode { ENFORCED, UNENFORCED }
@@ -442,10 +442,10 @@ public final class FleetMcp {
// fall back to here — AuthorizationMode governs enforcement, not identity.
Principal p = callers.resolve(req.getRemoteAddr(), req.getRemotePort(),
req.getHeader("Authorization"));
// CB-532: guard on the ROLE, not on the terminal being null. This excludes a
// lead, which carries its pane too, while including every spawned member role.
// Enrolling a lead would count it as an available member in the roster.
markSpawnedMemberPresent(p, presence);
// Guards on the ROLE, not on the terminal being null — this marks presence for
// a worker, an architect, or the unconfigured-pane floor, and excludes a lead
// or a collaborator even though each carries its own pane too.
markTrackedCallerPresent(p, presence);
return McpTransportContext.create(Map.of(
CALLER_TERMINAL, orEmpty(p.terminal()),
CALLER_PID, Long.toString(p.pid()),
@@ -459,12 +459,14 @@ public final class FleetMcp {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_send", req.arguments()),
str(req.arguments(), "sessionId"));
if (denied != null) return denied;
String caller = callerTerminal(exchange);
Principal caller = principal(exchange);
String callerTerminal = caller.terminal();
String callerOwner = caller.ownerKey();
// CB-548: only a PRIMARY caller may claim the legacy singleton "primary" fallback.
// An architect delegates as its own pane but must never become the fallback that
// no-delegation inbox nudges target as if it were the primary (the per-target
// delegation map does not cure the singleton).
recordPrimarySingleton(primaryRegistry, caller, principal(exchange));
recordPrimarySingleton(primaryRegistry, callerTerminal, caller);
Map<String, Object> a = req.arguments();
String target = str(a, "sessionId");
String content = str(a, "content");
@@ -480,18 +482,18 @@ public final class FleetMcp {
// Answering a worker's fleet_ask (CB-205): resolve its blocked question and
// block for the worker's reply as it resumes the same turn. This is the same
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
return answer(messages, turnId, content, timeoutMs(a), caller);
return answer(messages, turnId, content, timeoutMs(a), callerOwner);
}
// CB-548: delegator ownership (which lead's reply nudge this worker routes to,
// CB-532) is recorded only once the send is ACCEPTED — MessageService has won the
// session lock and queued delivery — via the accepted-delivery callback, never at
// request time. A concurrent sender that times out BUSY therefore cannot steal a
// live turn's reply routing without ever owning the turn.
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, callerTerminal);
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), caller);
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), callerOwner);
};
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
// is "is this caller a worker at all", and it can only ever reply as itself.
@@ -514,7 +516,7 @@ public final class FleetMcp {
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null);
if (denied != null) return denied;
return status(messages, str(req.arguments(), "sessionId"), callerTerminal(exchange));
return status(messages, str(req.arguments(), "sessionId"), principal(exchange).ownerKey());
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> pollHandler =
(exchange, req) -> {
@@ -524,7 +526,8 @@ public final class FleetMcp {
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
if (denied != null) return denied;
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
return poll(messages, leadChannel, str(a, "ticket"), target, coordId,
principal(exchange).ownerKey());
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
// Acking removes a reply from the inbox, so it is a drain, not a read.
@@ -833,9 +836,14 @@ public final class FleetMcp {
return identity;
}
/** Mark a connected spawned member available for the injector readiness gate. */
static void markSpawnedMemberPresent(Principal caller, MemberPresence presence) {
if (caller.isSpawnedMember()) {
/**
* Mark a caller present for the injector readiness gate, when its deliverability depends on
* proving a live MCP contact: a worker, an architect, or the unconfigured-pane floor. A lead
* or a collaborator is excluded — each is already deliverable through its own named-registry
* entry.
*/
static void markTrackedCallerPresent(Principal caller, MemberPresence presence) {
if (caller.isSpawnedMember() || caller.isObserver()) {
presence.markPresent(caller.terminal());
}
}
@@ -891,8 +899,8 @@ public final class FleetMcp {
* fleetd #612 B3 — as {@link #quarantineSource()}, {@code public} for the same cross-package
* reason, for the real {@link LeadRollover} (or {@code null}) this daemon was assembled with.
* {@code FleetdLeadRolloverAssemblyTest} drives {@code open}/{@code confirm} on this exact
* instance and waits for the real continuation to send {@code /clear} and {@code bootstrapText}
* through the real {@code router.leadAgents()}.
* instance and waits for the real continuation to end the old pane, relaunch a fresh one, and
* send {@code bootstrapText} through the real {@code router.leadAgents()}.
*/
public LeadRollover leadRollover() {
return leadRollover;
@@ -909,7 +917,7 @@ public final class FleetMcp {
*/
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
Long timeoutMs, Runnable onAccepted, Set<String> profiles,
String callerTerminal) {
String callerOwner) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -919,7 +927,7 @@ public final class FleetMcp {
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
try {
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerTerminal), timeout);
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerOwner), timeout);
} catch (HerdrException e) {
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
}
@@ -929,15 +937,15 @@ public final class FleetMcp {
* {@code fleet_send} carrying a {@code turnId}: the primary's answer to a worker's
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
* it resumes the same turn — surfaced to the primary identically to a normal send.
* {@code callerTerminal} must match the turn's recorded owner or this is refused.
* {@code callerOwner} must match the turn's recorded owner or this is refused.
*/
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs,
String callerTerminal) {
String callerOwner) {
if (isBlank(turnId) || isBlank(content)) {
return error("turnId and content are required to answer a worker's question");
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
return formatReply(messages.answer(turnId, content, timeout, callerTerminal), timeout);
return formatReply(messages.answer(turnId, content, timeout, callerOwner), timeout);
}
/**
@@ -1013,13 +1021,12 @@ public final class FleetMcp {
}
/**
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
* result back to that same caller — see {@link MessageService#poll(String, String)}.
* As above, recording {@code creator}'s owner key so a later {@code fleet_poll{ticket}} only
* hands the result back to the same caller — see {@link MessageService#poll(String, String)}.
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles,
String creatorTerminal) {
Runnable onAccepted, Set<String> profiles,
Principal creator) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -1027,7 +1034,7 @@ public final class FleetMcp {
if (targetError != null) {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
String ticket = messages.sendAsync(sessionId, content, onAccepted, creator);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
@@ -1217,13 +1224,12 @@ public final class FleetMcp {
}
/**
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
* client-supplied value.
* As above, refusing a ticket lookup whose caller owner key differs from the key that created it
* — see {@link MessageService#poll(String, String)}. {@code callerOwner} comes from the calling
* connection's resolved principal, never a client-supplied value.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId, String callerTerminal) {
String target, String coordId, String callerOwner) {
if (!isBlank(coordId)) {
return pollHeldPeerMail(leadChannel, coordId);
}
@@ -1237,7 +1243,7 @@ public final class FleetMcp {
if (isBlank(ticket)) {
return error("ticket (or target) is required");
}
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
MessageService.TaskView v = messages.poll(ticket, callerOwner);
if (v == null) {
return error("unknown ticket: " + ticket + " (never issued, or expired)");
}
@@ -1346,18 +1352,17 @@ public final class FleetMcp {
* {@code fleet_status}: the live lifecycle status of a worker session, plus — when the worker
* is paused mid-turn in an async {@code fleet_ask} — the open question and how to answer it, so
* a lead on its normal poll cadence does not need the ticket to notice. The question, its
* {@code turnId} and its ticket id are shown only to the caller whose terminal created that
* delegation, or to a caller with no terminal at all (the unnamed primary); any other caller
* still sees the base status. {@code callerTerminal} is the CALLING session's terminal id,
* resolved by the MCP layer from the connection, never a client-supplied value.
* {@code turnId} and its ticket id are shown only to the caller whose owner key created that
* delegation, or to the unnamed primary; any other caller still sees the base status.
* {@code callerOwner} comes from the calling connection's resolved principal.
*/
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerTerminal) {
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerOwner) {
if (isBlank(sessionId)) {
return error("sessionId is required");
}
try {
String base = messages.status(sessionId).name().toLowerCase();
MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerTerminal);
MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerOwner);
if (ask == null) {
return text(base);
}
@@ -1412,6 +1417,14 @@ public final class FleetMcp {
}
return text(json(m));
}
if (caller.isObserver()) {
// No architect slot, collaborator name, or lead name to report — only the pane itself,
// so a peer that already knows this terminal can still address it.
if (caller.terminal() != null) {
m.put("sessionId", caller.terminal());
}
return text(json(m));
}
if (!caller.isWorker()) {
// CB-530: which lead, once more than one pane is configured as one. `role` deliberately
// still reads "primary" — the fallback ladder in CLAUDE.md keys on it, and a lead IS a
@@ -1492,7 +1505,7 @@ public final class FleetMcp {
}
if (isBlank(callerTerminal)) {
// An unnamed primary (token/loopback path, no resolved pane) has nowhere for the
// eventual /clear + bootstrap to land — LeadRollover#open would throw
// eventual relaunch + bootstrap to land — LeadRollover#open would throw
// IllegalArgumentException for the same reason; refuse cleanly here instead.
return error("fleet_handover requires a named lead pane (a resolved connection terminal) "
+ "to open a rollover request against — an unnamed primary has none");
@@ -2630,17 +2643,20 @@ public final class FleetMcp {
private static McpSchema.Tool handoverTool() {
return tool(FleetTool.HANDOVER.wireName(),
"Replace your OWN lead session once its context is full: write a handover file, "
+ "then use this to have fleetd clear your pane and bootstrap a fresh lead "
+ "session against it. Four actions: 'open' (requests a token and the "
+ "handoverPath you must write the handover file to before confirming), "
+ "'confirm' (validates every gate and — only if every one passes — schedules "
+ "the roll; it does NOT itself clear the pane, the roll runs once this call's "
+ "own turn ends), 'cancel' (drops a pending request without rolling), and "
+ "'status' (read-only: what happened to a token after 'confirm' — still "
+ "running (approved but not finished yet), the roll completed, the calling "
+ "turn never settled within turnSettleSeconds so no /clear was ever sent, or "
+ "/clear itself never settled so bootstrapText was never sent; never "
+ "schedules, cancels or retries anything). Primary-only. "
+ "then use this to have fleetd end your pane's process and relaunch a fresh "
+ "lead session bootstrapped against it. Four actions: 'open' (requests a "
+ "token and the handoverPath you must write the handover file to before "
+ "confirming), 'confirm' (validates every gate and — only if every one "
+ "passes — schedules the roll; it does NOT itself end your pane, the roll "
+ "runs once this call's own turn ends), 'cancel' (drops a pending request "
+ "without rolling), and 'status' (read-only: what happened to a token after "
+ "'confirm' — still running (approved but not finished yet), the roll "
+ "completed, the calling turn never settled within turnSettleSeconds so "
+ "nothing was touched, the old pane never confirmed dead so no relaunch was "
+ "attempted, the relaunch itself failed, the fresh pane never became ready "
+ "so bootstrapText was never sent, or the fresh pane became ready and was "
+ "bootstrapped but was never recognised as a live lead; never schedules, "
+ "cancels or retries anything). Primary-only. "
+ "There is deliberately no terminal/session/leadTerminal parameter: the pane "
+ "to roll is always resolved from YOUR OWN connection, never a value you "
+ "pass, so you can only ever roll yourself — never another lead. Requires "
@@ -1,5 +1,6 @@
package dev.ltms.fleet.msg;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrRouter;
@@ -282,18 +283,15 @@ public final class MessageService {
*/
private volatile boolean askTimedOut;
/**
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
* created through an overload that does not record one. {@link #poll(String, String)}
* compares a polling caller's own terminal against this field before handing back the
* ticket's state.
* The owner key of the caller whose {@code fleet_send{wait:false}} created this ticket, or
* {@code null} for the unnamed primary and overloads that do not record a caller.
*/
private final String creatorTerminal;
private final String creatorOwner;
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
private Task(String ticket, String target, LongSupplier nowNanos, String creatorOwner) {
this.ticket = ticket;
this.target = target;
this.creatorTerminal = creatorTerminal;
this.creatorOwner = creatorOwner;
this.createdNanos = nowNanos.getAsLong();
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
@@ -355,6 +353,14 @@ public final class MessageService {
*/
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
/**
* Minted once per {@code MessageService} instance and folded into every ticket id (see
* {@link #sendAsync(String, String, Runnable, Principal)}). {@link #ticketSeq} alone restarts at
* zero for every instance, so without this a ticket id can be reused across instances and
* resolve to an unrelated {@link Task} with no error; this nonce makes that impossible, because
* an id minted by one instance can never match the id space of another.
*/
private final String ticketBootNonce = UUID.randomUUID().toString().substring(0, 6);
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual().name("bridge-async-", 0).factory());
@@ -930,13 +936,13 @@ public final class MessageService {
/**
* Deliver {@code content} to {@code target} (a herdr {@code terminal_id}) and block until the
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. {@code callerTerminal}
* is the terminal of the caller making this call — {@code null} for the unnamed primary — and is
* recorded as the turn's owner, the only caller {@link #answer(String, String, long, String)} will
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. {@code callerOwner}
* identifies the caller making this call and is recorded as the turn's owner. It is the only
* caller {@link #answer(String, String, long, String)} will
* later accept an answer from if the worker pauses mid-turn to ask.
*/
public Reply send(String target, String content, long timeoutMillis, String callerTerminal) {
return send(target, content, timeoutMillis, null, callerTerminal);
public Reply send(String target, String content, long timeoutMillis, String callerOwner) {
return send(target, content, timeoutMillis, null, callerOwner);
}
/**
@@ -951,13 +957,13 @@ public final class MessageService {
* acceptance means a concurrent sender that times out {@code BUSY} can never steal ownership it
* never earned. {@code null} disables the hook.
*/
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerTerminal) {
return send(target, content, timeoutMillis, onAccepted, null, callerTerminal);
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerOwner) {
return send(target, content, timeoutMillis, onAccepted, null, callerOwner);
}
/** Run a send, optionally stopping an async task that teardown already failed before acceptance. */
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task,
String callerTerminal) {
String callerOwner) {
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
@@ -977,7 +983,7 @@ public final class MessageService {
// open race). Opening first also means a throwing onAccepted (fired before enqueue) or an
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
// failed send leaves no stale waiter behind.
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target, Rendezvous.Owner.of(callerTerminal));
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target, Rendezvous.Owner.of(callerOwner));
// CB-640: this send now owns target's delivery, so any earlier stranded-reply or
// still-queued fact no longer describes the live state — clear both rather than let
// them outlive the send that supersedes them.
@@ -1157,20 +1163,20 @@ public final class MessageService {
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
* is unblocked so a reply that lands the instant it resumes is not lost.
*
* <p>{@code callerTerminal} is the terminal of the caller making this call — {@code null} for
* the unnamed primary. It is checked against the turn's recorded owner (the caller whose
* <p>{@code callerOwner} identifies the caller making this call. It is checked against the
* turn's recorded owner (the caller whose
* accepted delegation opened it, see {@link #send(String, String, long, String)} and
* {@link #sendAsync(String, String, Runnable, String)}) before anything else runs: a mismatch,
* {@link #sendAsync(String, String, Runnable, Principal)}) before anything else runs: a mismatch,
* including a turn with no owner on record at all, returns {@link Outcome#NOT_TURN_OWNER}
* without touching the rendezvous, the session lock, or any async task bookkeeping.
*/
public Reply answer(String turnId, String content, long timeoutMillis, String callerTerminal) {
public Reply answer(String turnId, String content, long timeoutMillis, String callerOwner) {
String workerSession = rendezvous.askSession(turnId);
if (workerSession == null) {
return new Reply(Outcome.STALE_TURN, null); // the ask lapsed (timed out or already answered)
}
Rendezvous.Owner owner = rendezvous.askOwner(turnId);
if (!Rendezvous.Owner.permits(owner, callerTerminal)) {
if (!Rendezvous.Owner.permits(owner, callerOwner)) {
return new Reply(Outcome.NOT_TURN_OWNER, null);
}
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
@@ -1313,16 +1319,16 @@ public final class MessageService {
}
/**
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
* differs from this one; {@code null} records no owner (a caller with no terminal — the
* unnamed primary — is always allowed to poll the result regardless).
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creator}'s owner key as this
* ticket's owner. The key is derived here from the resolved principal so callers cannot pass a
* terminal address where an owner identity is required.
*
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
String ticket = "task-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
public String sendAsync(String target, String content, Runnable onAccepted, Principal creator) {
String ticket = "task-" + ticketBootNonce + "-" + ticketSeq.incrementAndGet();
String creatorOwner = creator == null ? null : creator.ownerKey();
Task task = new Task(ticket, target, nowNanos, creatorOwner);
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
@@ -1353,7 +1359,7 @@ public final class MessageService {
}
asyncExecutor.submit(() -> {
try {
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorTerminal);
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorOwner);
if (result.outcome() == Outcome.QUESTION) {
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
// just after resolveQuestion wakes this thread.
@@ -1385,9 +1391,9 @@ public final class MessageService {
}
/**
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
* As {@link #poll(String, String)}, with no caller owner key — the unnamed primary's ticket rule
* checked, so this overload must only be used where the caller's identity is otherwise
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
* irrelevant.
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
@@ -1395,18 +1401,18 @@ public final class MessageService {
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
* Refuses a {@code callerOwner} that differs from the owner that created the ticket (see
* {@link #sendAsync(String, String, Runnable, Principal)}) with a {@link Phase#FAILED} view that
* carries no reply text. The unnamed primary has a {@code null} owner key and is never refused.
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
*/
public TaskView poll(String ticket, String callerTerminal) {
public TaskView poll(String ticket, String callerOwner) {
Task task = tasks.get(ticket);
if (task == null) {
return null;
}
if (!ownsTicket(task, callerTerminal)) {
if (!ownsTicket(task, callerOwner)) {
return new TaskView(ticket, Phase.FAILED, null, null,
"forbidden: this ticket was created by a different session", null);
}
@@ -1446,14 +1452,13 @@ public final class MessageService {
}
/**
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
* equal the terminal recorded on the task; a task with no recorded terminal matches no
* terminal-bearing caller.
* Whether {@code callerOwner} may read {@code task}'s state. A {@code null} caller key is the
* unnamed primary and may read every ticket. Other callers must match the task's owner key. This
* differs from {@link Rendezvous.Owner#permits}: a missing rendezvous owner is not an authenticated
* unnamed primary, so that gate refuses every caller when no owner was recorded.
*/
private static boolean ownsTicket(Task task, String callerTerminal) {
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
private static boolean ownsTicket(Task task, String callerOwner) {
return callerOwner == null || callerOwner.equals(task.creatorOwner);
}
/**
@@ -1779,13 +1784,13 @@ public final class MessageService {
* {@code fleet_status} uses this to show a pending question without the caller needing the
* ticket. {@code null} when the session has no open async question (including a session mid a
* <em>blocking</em> {@code fleet_ask}, which has no {@link Task} to look up — see
* {@link PendingAsk}), or when {@code callerTerminal} does not own the task the question
* {@link PendingAsk}), or when {@code callerOwner} does not own the task the question
* belongs to (see {@link #ownsTicket(Task, String)}).
*/
public PendingAsk pendingAsk(String workerSession, String callerTerminal) {
public PendingAsk pendingAsk(String workerSession, String callerOwner) {
for (Task task : tasks.values()) {
Reply q = task.question;
if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerTerminal)) {
if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerOwner)) {
return new PendingAsk(task.ticket, q.text(), q.turnId());
}
}
@@ -1,5 +1,6 @@
package dev.ltms.fleet.msg;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
@@ -72,23 +73,22 @@ public final class Rendezvous {
/**
* The caller whose accepted delegation opened a turn — the only caller allowed to answer it.
* A {@code null} terminal means the unnamed primary, an authenticated caller with no pane.
* A {@code null} owner key means the unnamed primary.
*/
public record Owner(String terminal) {
public record Owner(String ownerKey) {
public static final Owner UNNAMED_PRIMARY = new Owner(null);
public static Owner of(String terminal) {
return terminal == null ? UNNAMED_PRIMARY : new Owner(terminal);
public static Owner of(String ownerKey) {
return ownerKey == null ? UNNAMED_PRIMARY : new Owner(ownerKey);
}
/**
* Whether {@code callerTerminal} matches {@code owner}. A {@code null} owner means no
* owner was ever recorded, and that state matches no caller, not even one whose own
* terminal is {@code null} — "no record" and "recorded as the unnamed primary" are
* different states.
* Whether {@code callerOwner} matches {@code owner}. A {@code null} owner means no owner was
* recorded, so it matches no caller. {@link #UNNAMED_PRIMARY} records the unnamed primary
* with an owner object whose key is {@code null}.
*/
public static boolean permits(Owner owner, String callerTerminal) {
return owner != null && java.util.Objects.equals(owner.terminal(), callerTerminal);
public static boolean permits(Owner owner, String callerOwner) {
return owner != null && java.util.Objects.equals(owner.ownerKey(), callerOwner);
}
}
@@ -101,6 +101,14 @@ public final class Rendezvous {
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */
private final ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
private final AtomicLong askSeq = new AtomicLong();
/**
* Minted once per {@code Rendezvous} instance and folded into every {@code turnId} (see
* {@link #openAsk(String)}). {@link #askSeq} alone restarts at zero for every instance, so
* without this a {@code turnId} minted by one instance could be minted again by another and
* resolve to an unrelated ask with no error; this nonce makes that impossible, because an id
* minted by one instance can never match the id space of another.
*/
private final String askBootNonce = UUID.randomUUID().toString().substring(0, 6);
/** Per-session index of the currently-open ask, so duplicate fleet_ask calls coalesce onto one turn. */
private final ConcurrentHashMap<String, String> openAsksBySession = new ConcurrentHashMap<>();
@@ -187,7 +195,7 @@ public final class Rendezvous {
while (true) {
AskWaiter[] minted = { null };
String turnId = openAsksBySession.computeIfAbsent(session, _ -> {
String newTurnId = session + "#" + askSeq.incrementAndGet();
String newTurnId = session + "#" + askBootNonce + "-" + askSeq.incrementAndGet();
CompletableFuture<String> answer = new CompletableFuture<>();
AskWaiter waiter = new AskWaiter(session, answer, ownerOf(session));
asks.put(newTurnId, waiter);
@@ -664,7 +664,7 @@ public final class FleetApp {
return;
}
Principal caller = ctx.attribute(CALLER);
String callerTerminal = caller == null ? null : caller.terminal();
String callerOwner = caller == null ? null : caller.ownerKey();
JsonNode body;
try {
body = mapper.readTree(ctx.body());
@@ -690,19 +690,19 @@ public final class FleetApp {
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
if (turnId != null && !turnId.isBlank()) {
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerTerminal), timeout);
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerOwner), timeout);
return;
}
if (!wait) {
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
String ticket = messages.sendAsync(id, content, null, callerTerminal);
String ticket = messages.sendAsync(id, content, null, caller);
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
return;
}
try {
writeReply(ctx, id, messages.send(id, content, timeout, callerTerminal), timeout);
writeReply(ctx, id, messages.send(id, content, timeout, callerOwner), timeout);
} catch (HerdrException e) {
herdrError(ctx, e);
}
@@ -883,10 +883,10 @@ public final class FleetApp {
body.put("ready", deliverable.test(id));
// A worker paused mid-turn in an async fleet_ask is otherwise invisible to a status
// poll — surface the open question and how to answer it, same as fleet_poll's
// Phase.ASKING view, but only to the caller whose terminal created that delegation, or
// to a caller with no terminal at all (the unnamed primary).
// Phase.ASKING view, but only to the caller whose owner key created that delegation, or
// to the unnamed primary.
Principal caller = ctx.attribute(CALLER);
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.terminal());
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.ownerKey());
if (ask != null) {
body.put("question", ask.question());
body.put("turnId", ask.turnId());
@@ -904,7 +904,7 @@ public final class FleetApp {
return;
}
Principal caller = ctx.attribute(CALLER);
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.terminal());
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.ownerKey());
if (v == null) {
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
return;
@@ -255,6 +255,10 @@ public final class SessionManager implements TurnListener {
handle.id(), handle.terminalId(), resolvedProfile, actualRole, cwd, ownerTerminal, now, now, 0,
MemberSession.State.SPAWNING, null, null, handle.charterReceipt(), handle.agentSessionId());
registry.put(handle.id(), session);
// A presence contact that already arrived for this terminal found no registry
// entry to transition and gave up silently. Retry it now that one exists; remove
// this call and such a session stays in SPAWNING even though it is present.
reconcilePresence(handle.terminalId());
handles.put(handle.id(), handle);
log.debug("acquired session id={} terminal={} profile={} owner={}",
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
@@ -466,6 +470,12 @@ public final class SessionManager implements TurnListener {
MemberSession resolved = resolveAgentSessionId(removed, removedHandle);
notifyReleased(new ReleaseDetail(resolved.terminalId(), resolved.worktree(),
resolved.branch(), snapshotRef, resolved.agentSessionId()));
String terminal = removed.terminalId();
if (terminal != null && !terminal.isBlank()) {
// Without this, a terminal stays marked present after its pane is gone, so a
// later send to the same id would read as deliverable instead of refused.
presence.forget(terminal);
}
}
}
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
@@ -809,6 +819,10 @@ public final class SessionManager implements TurnListener {
handle.charterReceipt(),
handle.agentSessionId());
registry.put(handle.id(), session);
// A presence contact that already arrived for this terminal found no registry entry to
// transition and gave up silently. Retry it now that one exists; remove this call and
// such a session stays in SPAWNING even though it is present.
reconcilePresence(handle.terminalId());
handles.put(handle.id(), handle);
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
@@ -983,6 +997,20 @@ public final class SessionManager implements TurnListener {
transitionByTerminal(terminalId, MemberSession.State.SPAWNING, MemberSession.State.READY);
}
/**
* Completes a newly registered session's {@code SPAWNING -> READY} transition when {@code
* terminalId} was already marked present before this ran. A terminal never marked present is
* left in {@code SPAWNING}; it reaches {@code READY} normally through {@link #onReady} once
* its own contact arrives. Callers must run this only once the session's registry entry is
* already visible — {@link #onReady}'s transition matches against that entry, and reconciling
* before the entry exists finds nothing to transition.
*/
private void reconcilePresence(String terminalId) {
if (terminalId != null && !terminalId.isBlank() && presence.isPresent(terminalId)) {
onReady(terminalId);
}
}
/**
* Lifecycle hook: a message was delivered into the worker — it is now busy on a turn.
* The turn count is bumped and the activity timestamp is refreshed. A {@code DONE} session
@@ -6,6 +6,8 @@ import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.lead.LeadRollover;
import dev.ltms.fleet.msg.ReplyInbox;
import io.javalin.Javalin;
@@ -34,7 +36,7 @@ import static org.junit.jupiter.api.Assertions.fail;
* fleetd #612 Unit A) with three methods: {@code unrelatedAnchorStillPresent} (a scaffold anchor,
* not an independent claim — needs no replacement of its own), {@code
* mainStillCallsTheLeadRolloverFactory} (the call-site pin replaced by {@link
* #assembledLeadRolloverRunsTheRealClearAndBootstrapSequence}), and {@code
* #assembledLeadRolloverEndsTheOldPaneThroughTheRealHerdrRouter}), and {@code
* factoryGatesOnConfigPresence} (the absent-config claim replaced by {@link
* #absentLeadRolloverConfigMeansNoRolloverIsBuilt} — a claim this ticket found was NOT actually
* covered behaviourally anywhere else: {@code LeadRolloverTest}'s only related assertion is
@@ -50,8 +52,8 @@ import static org.junit.jupiter.api.Assertions.fail;
* invisible to this test, even though the two are genuinely different daemons in production. This
* version configures two distinct sockets and two distinct {@link FakeHerdr} instances (the same
* pattern {@code FleetdAssemblyConnectionIdentityTest}, fleetd #612 B2, already uses to separate
* lead from member) and asserts the roll's {@code /clear}/bootstrap sends land on the LEAD fake
* and never on the MEMBER one.
* lead from member) and asserts the roll's {@code pane.close} call lands on the LEAD fake and
* never on the MEMBER one.
*/
class FleetdLeadRolloverAssemblyTest {
@@ -177,11 +179,47 @@ class FleetdLeadRolloverAssemblyTest {
return FleetConfig.load(f);
}
@SuppressWarnings("unchecked")
/**
* Unlike {@link #writeConfig}, this names a {@code profile:} for the lead and declares it
* under {@code profiles:}, so {@code LeadLauncher#relaunch} can actually start a fresh agent
* instead of refusing with "names no profile". {@code relaunchReadySeconds} is cut to 2s so
* the recognition wait (expected to time out — see the test) does not cost real test seconds.
*/
private static FleetConfig writeConfigWithRelaunchableLead(Path dir, Path leadCwd) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: "%s"
memberHerdrSocket: "%s"
idleSleepGuard:
enabled: false
broker:
uri: "amqp://fake-test-broker/vh"
fleet:
leaders:
opus:
tab: "lead: opus"
cwd: "%s"
profile: opus
profiles:
opus:
subscription: true
argv: ["ccs", "opus"]
leadRollover:
handoverPath: handover.md
requireOperatorConfirm: false
relaunchReadySeconds: 2
""".formatted(LEAD_SOCKET, MEMBER_SOCKET, leadCwd.toString()));
return FleetConfig.load(f);
}
@Test
@DisplayName("[BEHAVIOURAL] the real assembled LeadRollover runs the full open/confirm/continuation "
+ "sequence — /clear, then bootstrapText — through the real herdr router")
void assembledLeadRolloverRunsTheRealClearAndBootstrapSequence(@TempDir Path dir) throws Exception {
@DisplayName("[BEHAVIOURAL] the real assembled LeadRollover runs the open/confirm/continuation "
+ "sequence through the real herdr router — ending the old pane, then giving up once it "
+ "never reports gone")
void assembledLeadRolloverEndsTheOldPaneThroughTheRealHerdrRouter(@TempDir Path dir) throws Exception {
Path leadCwd = dir.resolve("lead-workspace");
Files.createDirectories(leadCwd);
FleetConfig cfg = writeConfig(dir, leadCwd);
@@ -204,6 +242,11 @@ class FleetdLeadRolloverAssemblyTest {
+ "Fleetd.leadRollover(...) call site — a mutation to `LeadRollover leadRollover = "
+ "null;` at that call site can never pass this");
// assembleAndStart's own boot work (the orphan-worker reap) makes a real call on the
// member daemon before the roll ever starts. Clear it here so the assertion below measures
// only what the roll itself does, not what daemon startup does.
member.calls.clear();
LeadRollover.PendingRollover pending = rollover.open("term_a", "fleetd #612 B3 test");
String expectedHandoverPath = leadCwd.resolve("handover.md").normalize().toString();
assertEquals(expectedHandoverPath, pending.handoverPath());
@@ -223,40 +266,109 @@ class FleetdLeadRolloverAssemblyTest {
// constructor), so this polls the real FleetMcp.leadRollover() instance's status(token)
// until the real continuation finishes.
LeadRollover.RollStatus status = pollUntilTerminal(rollover, pending.token());
assertEquals(LeadRollover.RollState.ROLLED, status.state(),
"the full happy path must complete: FakeHerdr's default agent status is 'idle', so "
+ "the turn-boundary wait settles immediately and the post-/clear wait "
+ "releases via its pickup-grace path — detail: " + status.detail());
// Prove the real herdr router actually sent BOTH messages, in order, to the real LEAD
// pane — this is the one thing a source-text pin on the call site could never show.
// FakeHerdr's pane.get is a fixed canned response that never reports a pane as gone, so the
// real router's death poll runs out its whole budget and the roll stops here — proving the
// real teardown call landed on the real LEAD pane without ever reaching a relaunch or a send.
assertEquals(LeadRollover.RollState.OLD_PANE_NEVER_DIED, status.state(),
"the old pane never reports gone against this fake, so the roll must stop with "
+ "OLD_PANE_NEVER_DIED rather than ever relaunching or sending anything — "
+ "detail: " + status.detail());
// Prove the real herdr router actually closed the real LEAD pane — this is the one thing a
// source-text pin on the call site could never show.
boolean closedOldPane = lead.calls.stream()
.anyMatch(c -> c.method().equals("pane.close")
&& c.params() instanceof Map<?, ?> m && "w2:p7".equals(m.get("pane_id")));
assertTrue(closedOldPane, "endOldSession must close the real old pane (w2:p7) through the "
+ "real LEAD herdr client, got calls: " + lead.calls);
// No agent.prompt is ever sent on this path: the roll stops at the pane-death wait, strictly
// before the relaunch and the final send step.
List<FakeHerdr.Call> prompts = lead.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.toList();
assertTrue(prompts.size() >= 2, "expected at least a /clear send and a bootstrapText send "
+ "on the LEAD daemon, got " + prompts.size() + " agent.prompt calls: " + prompts);
assertEquals("/clear", ((Map<String, Object>) prompts.get(0).params()).get("text"),
"the first send must be the literal /clear housekeeping command");
Object secondText = ((Map<String, Object>) prompts.get(1).params()).get("text");
assertTrue(secondText instanceof String && ((String) secondText).contains(expectedHandoverPath),
"the second send must be the default bootstrapText naming the resolved handover "
+ "path, got: " + secondText);
assertTrue(prompts.isEmpty(), "a roll that stops at OLD_PANE_NEVER_DIED must never reach the "
+ "send step, got agent.prompt call(s) on the LEAD daemon: " + prompts);
// fleetd #612 B3 correction: prove the roll never touches the MEMBER daemon. A mutation
// swapping router.leadAgents() for router.memberAgents() at the real call site would move
// both sends above onto `member` instead, which this assertion catches — the thing the
// single-fake version of this test could never see, because both wrapped the same client.
// the pane.close call above onto `member` instead, which this assertion catches — the thing
// the single-fake version of this test could never see, because both wrapped the same client.
assertTrue(member.calls.isEmpty(), "the roll must be wired to the LEAD daemon only — got "
+ member.calls.size() + " call(s) recorded on the MEMBER daemon since the roll began: "
+ member.calls);
}
/**
* Exercises the relaunch site {@link #assembledLeadRolloverEndsTheOldPaneThroughTheRealHerdrRouter}
* never reaches: with the old pane confirmed gone, the roll relaunches a fresh lead, and
* {@code bootstrapText} must reach it even though recognition times out (FakeHerdr's
* {@code tab.list} is a fixed canned response that never reflects the relaunch's own
* {@code tab.rename}, so the fresh terminal is never recognised as a live lead). Same
* dual-socket shape as the sibling test: two distinct {@link FakeHerdr} instances, so a
* {@code bootstrapText} send wired to the wrong daemon is visible.
*/
@Test
@DisplayName("[BEHAVIOURAL] bootstrapText reaches the fresh LEAD terminal even when recognition "
+ "times out, and the MEMBER daemon never sees it")
void bootstrapTextReachesTheFreshLeadTerminalEvenWhenRecognitionTimesOut(@TempDir Path dir) throws Exception {
Path leadCwd = dir.resolve("lead-workspace");
Files.createDirectories(leadCwd);
FleetConfig cfg = writeConfigWithRelaunchableLead(dir, leadCwd);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
RecordingResourcePorts ports = new RecordingResourcePorts();
FakeHerdr lead = new FakeHerdr();
lead.withTab("w2", "w2:t7", "lead: opus");
// Lets the old pane (w2:p7) report gone once pane.close actually reaches it, so the roll
// proceeds to relaunch instead of stopping at OLD_PANE_NEVER_DIED.
lead.paneGoneAfterClose("w2:p7");
FakeHerdr member = new FakeHerdr();
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
LeadRollover rollover = runtime.mcp().leadRollover();
assertNotNull(rollover, "leadRollover: is present in this test's config, so a real "
+ "LeadRollover must have been built");
LeadRollover.PendingRollover pending = rollover.open("term_a", "bootstrapText relaunch test");
Thread.sleep(50);
Files.writeString(Path.of(pending.handoverPath()), "handover content for bootstrapText test");
LeadRollover.RollDecision decision = rollover.confirm("term_a", pending.token(), true);
assertTrue(decision.accepted(), "confirm() must approve — got: " + decision);
LeadRollover.RollStatus status = pollUntilTerminal(rollover, pending.token());
assertEquals(LeadRollover.RollState.RELAUNCH_NOT_RECOGNISED, status.state(),
"the fresh terminal is never recognised against this fake's static tab.list, so the "
+ "roll must reach RELAUNCH_NOT_RECOGNISED — not an earlier failure state and "
+ "not ROLLED — detail: " + status.detail());
@SuppressWarnings("unchecked")
List<FakeHerdr.Call> leadPrompts = lead.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.toList();
assertEquals(1, leadPrompts.size(), "exactly one bootstrapText send is expected, on the LEAD "
+ "daemon, once recognition gives up — got: " + leadPrompts);
Object text = ((Map<String, Object>) leadPrompts.get(0).params()).get("text");
assertTrue(text instanceof String && ((String) text).contains("handover.md"),
"the send must be bootstrapText naming the resolved handover path, got: " + text);
// Scoped to agent.prompt specifically, not every MEMBER call: the orphan-worker reap also
// talks to the MEMBER daemon once, unconditionally, at daemon boot — unrelated to this roll.
List<FakeHerdr.Call> memberPrompts = member.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.toList();
assertTrue(memberPrompts.isEmpty(), "the roll must be wired to the LEAD daemon only — got "
+ memberPrompts.size() + " agent.prompt call(s) on the MEMBER daemon instead: "
+ memberPrompts);
assertTrue(memberPrompts.isEmpty(), "bootstrapText must never be sent to the MEMBER daemon, "
+ "got: " + memberPrompts);
}
private static LeadRollover.RollStatus pollUntilTerminal(LeadRollover rollover, String token)
throws InterruptedException {
long deadline = System.nanoTime() + java.util.concurrent.TimeUnit.SECONDS.toNanos(10);
long deadline = System.nanoTime() + java.util.concurrent.TimeUnit.SECONDS.toNanos(15);
while (System.nanoTime() < deadline) {
LeadRollover.RollStatus status = rollover.status(token);
if (status.state() != LeadRollover.RollState.PENDING
@@ -283,8 +395,10 @@ class FleetdLeadRolloverAssemblyTest {
""");
ConfigRef config = new ConfigRef(yaml, FleetConfig.load(yaml));
AgentControl agents = new AgentControl(new FakeHerdr());
WorkspaceControl spaces = new WorkspaceControl(new FakeHerdr());
LeadLauncher launcher = new LeadLauncher(agents, spaces, config.get());
LeadRollover rollover = Fleetd.leadRollover(config.get(), agents, config, Map::of);
LeadRollover rollover = Fleetd.leadRollover(config.get(), agents, spaces, launcher, config, Map::of);
assertNull(rollover, "leadRollover: is absent from this config, so the factory's opt-in "
+ "gate (`if (cfg.leadRollover() == null) return null;`) must fire and no "
@@ -4,6 +4,8 @@ import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.lead.LeadRollover;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
@@ -54,6 +56,11 @@ class FleetdLeadRolloverWorkspaceLookupTest {
return new AgentControl(new FakeHerdr());
}
/** None of this class's tests reach the deferred continuation, so a plain fake is enough. */
private static LeadLauncher fakeLauncher(FleetConfig cfg) {
return new LeadLauncher(fakeAgents(), new WorkspaceControl(new FakeHerdr()), cfg);
}
@Test
@DisplayName("[BEHAVIOURAL] Fleetd.leadRollover(...) resolves a relative handoverPath against "
+ "the CALLING lead's configured cwd, not the daemon's own working directory")
@@ -74,7 +81,8 @@ class FleetdLeadRolloverWorkspaceLookupTest {
""".formatted(leadCwd.toString()));
ConfigRef config = new ConfigRef(yaml, FleetConfig.load(yaml));
LeadRollover rollover = Fleetd.leadRollover(config.get(), fakeAgents(), config,
LeadRollover rollover = Fleetd.leadRollover(config.get(), fakeAgents(),
new WorkspaceControl(new FakeHerdr()), fakeLauncher(config.get()), config,
() -> Map.of("term_opus", "opus"));
assertNotNull(rollover, "leadRollover: is present in the loaded config, so the factory "
+ "must construct an object");
@@ -105,7 +113,8 @@ class FleetdLeadRolloverWorkspaceLookupTest {
// No lead has been discovered yet — exactly the real shape of a lead the live tab scan
// has not yet scanned, or one with no fleet.leaders entry at all.
LeadRollover rollover = Fleetd.leadRollover(config.get(), fakeAgents(), config, Map::of);
LeadRollover rollover = Fleetd.leadRollover(config.get(), fakeAgents(),
new WorkspaceControl(new FakeHerdr()), fakeLauncher(config.get()), config, Map::of);
assertNotNull(rollover);
LeadRollover.PendingRollover pending = rollover.open("term_unknown", "test");
@@ -144,7 +153,8 @@ class FleetdLeadRolloverWorkspaceLookupTest {
// below — exactly the natural mistake to make, since leads are discovered by a live tab
// scan that runs AFTER this factory is constructed at startup.
Map<String, String> liveLeadTerminals = new HashMap<>();
LeadRollover rollover = Fleetd.leadRollover(config.get(), fakeAgents(), config,
LeadRollover rollover = Fleetd.leadRollover(config.get(), fakeAgents(),
new WorkspaceControl(new FakeHerdr()), fakeLauncher(config.get()), config,
() -> liveLeadTerminals);
assertNotNull(rollover);
@@ -15,6 +15,7 @@ class AuthzTest {
private static final Principal ARCH_DESIGN = Principal.architect("lead-designer", "term_design", 400);
private static final Principal ARCH_OTHER = Principal.architect("reviewer", "term_review", 500);
private static final Principal COLLABORATOR = Principal.collaborator("ops", "term_collab", 600);
private static final Principal OBSERVER = Principal.observer("term_observer", 700);
@Test
void anonymousIsAuthorizedForNothing() {
@@ -266,4 +267,49 @@ class AuthzTest {
assertTrue(WORKER_A.isSpawnedMember());
assertTrue(ARCH_DESIGN.isSpawnedMember());
}
// ── the observer matrix ─────────────────────────────────────────────────────────────────────
@Test
void anObserverMayReadAndScrapeMetrics() {
assertTrue(Authz.permits(OBSERVER, READ, null));
assertTrue(Authz.permits(OBSERVER, METRICS, null));
}
@Test
void anObserverMayReplyAndAskOnlyAsItsOwnPane() {
assertTrue(Authz.permits(OBSERVER, REPLY, "term_observer"), "its own pane is its own");
assertTrue(Authz.permits(OBSERVER, ASK, "term_observer"));
assertFalse(Authz.permits(OBSERVER, REPLY, "term_design"),
"an observer must not reply on another pane");
assertFalse(Authz.permits(OBSERVER, REPLY, null),
"an absent target must not pass the own-session rule");
}
/**
* Every action beyond READ/METRICS/REPLY/ASK, asserted denied for an observer — including
* {@code TASK_READ}, which is the entire point of this role: an unconfigured pane must not be
* able to poll a ticket or read another session's status.
*/
@Test
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAndAsk() {
for (Authz.Action a : Authz.Action.values()) {
if (a == READ || a == METRICS || a == REPLY || a == ASK) {
continue;
}
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
"an observer must not " + a + " even when the classifier accepts every target");
}
}
@Test
void anObserverIsNotCountedAsAnyOtherRole() {
assertFalse(OBSERVER.isPrimary());
assertFalse(OBSERVER.isWorker());
assertFalse(OBSERVER.isArchitect());
assertFalse(OBSERVER.isCollaborator());
assertFalse(OBSERVER.isSpawnedMember());
assertTrue(OBSERVER.isObserver());
}
}
@@ -57,17 +57,23 @@ class CallerResolverTest {
return members;
}
/**
* With no roster wired up at all (the simple constructor), a loopback pane that owns a herdr
* pane but is not recognised as a live spawned member lands on the {@link Role#OBSERVER} floor
* — unforgeable and never token-gated, exactly like a worker's own identity, because it comes
* from the same connection-derived pane mapping.
*/
@Test
void aLoopbackWorkerPaneResolvesToWorkerRegardlessOfAuthMode() {
void aLoopbackPaneWithNoLiveRosterResolvesToObserverRegardlessOfAuthMode() {
Principal underTrust = new CallerResolver(workerIdentity()).resolve("127.0.0.1", 42, null);
Principal underToken = new CallerResolver(workerIdentity(), true, "s3cret")
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, underTrust.role());
assertEquals(Role.OBSERVER, underTrust.role());
assertEquals("term_a", underTrust.terminal());
assertEquals(Role.WORKER, underToken.role(),
"worker identity is unforgeable and must never be token-gated — otherwise enabling "
+ "auth would lock the whole fleet out of fleet_reply");
assertEquals(Role.OBSERVER, underToken.role(),
"the floor is unforgeable and must never be token-gated — otherwise enabling auth "
+ "would lock every unconfigured pane out of even READ");
assertEquals("term_a", underToken.terminal());
}
@@ -91,20 +97,20 @@ class CallerResolverTest {
}
@Test
void otherPanesRemainWorkersWhenAPinIsSet() {
void otherPanesRemainAtTheFloorWhenAPinIsSet() {
Principal p = CallerResolver.pinnedTo(workerIdentity(), false, null, "term_someone_else")
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertEquals("term_a", p.terminal());
}
/** The pin is optional config, so an absent or whitespace one must change nothing at all. */
@Test
void aBlankPinLeavesWorkerResolutionUntouched() {
assertEquals(Role.WORKER,
void aBlankPinLeavesFloorResolutionUntouched() {
assertEquals(Role.OBSERVER,
CallerResolver.pinnedTo(workerIdentity(), false, null, " ").resolve("127.0.0.1", 42, null).role());
assertEquals(Role.WORKER,
assertEquals(Role.OBSERVER,
CallerResolver.pinnedTo(workerIdentity(), false, null, null).resolve("127.0.0.1", 42, null).role());
}
@@ -225,13 +231,13 @@ class CallerResolverTest {
* mid-scan teardown into a refusal — the real match is still found and resolves as a worker.
*/
@Test
void aHerdrErrorOnANonOwningPaneStillResolvesTheRealWorker() {
void aHerdrErrorOnANonOwningPaneStillResolvesTheRealPane() {
FakeHerdr vanishedElsewhere = new FakeHerdr().processInfoFailsForPane("w2:p9", "pane_not_found");
ConnectionIdentity id = new ConnectionIdentity(new PaneLocator(vanishedElsewhere), _ -> FakeHerdr.WORKER_PID);
Principal p = new CallerResolver(id).resolve("127.0.0.1", 55555, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertEquals("term_a", p.terminal());
}
@@ -275,12 +281,12 @@ class CallerResolverTest {
}
@Test
void aPaneAbsentFromTheRegistryIsStillAWorker() {
void aPaneAbsentFromTheRegistryFallsToTheObserverFloor() {
Principal p = new CallerResolver(workerIdentity(), false, null,
Map.of("term_elsewhere", "gpt-sol-5.6"))
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertEquals("term_a", p.terminal());
assertNull(p.name());
}
@@ -306,16 +312,16 @@ class CallerResolverTest {
}
@Test
void anEmptyRegistryLeavesEveryPaneAWorker() {
void anEmptyRegistryLeavesEveryPaneAtTheObserverFloor() {
Map<String, String> noLeads = null;
assertEquals(Role.WORKER,
assertEquals(Role.OBSERVER,
new CallerResolver(workerIdentity(), false, null, Map.of())
.resolve("127.0.0.1", 42, null).role());
assertEquals(Role.WORKER,
assertEquals(Role.OBSERVER,
new CallerResolver(workerIdentity(), false, null, noLeads)
.resolve("127.0.0.1", 42, null).role());
// CB-531: and the same for the live-registry form, whose supplier may also be absent.
assertEquals(Role.WORKER,
// And the same for the live-registry form, whose supplier may also be absent.
assertEquals(Role.OBSERVER,
CallerResolver.withLeads(workerIdentity(), false, null, null)
.resolve("127.0.0.1", 42, null).role());
}
@@ -388,7 +394,7 @@ class CallerResolverTest {
Map<String, String> live = new java.util.HashMap<>();
CallerResolver r = CallerResolver.withLeads(workerIdentity(), false, null, () -> live);
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
assertEquals(Role.OBSERVER, r.resolve("127.0.0.1", 42, null).role());
live.put("term_a", "gpt-sol-5.6"); // the scanner sees a newly-labelled tab
@@ -405,7 +411,7 @@ class CallerResolverTest {
mutable.put("term_a", "sneaky");
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
assertEquals(Role.OBSERVER, r.resolve("127.0.0.1", 42, null).role());
}
// ── CB-548: architect slots ─────────────────────────────────────────────────────────────────
@@ -435,13 +441,13 @@ class CallerResolverTest {
}
@Test
void anUnboundPaneStillResolvesAsAWorker() {
void anUnboundPaneResolvesToTheObserverFloor() {
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(Map.of(),
Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null));
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, members).resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertNull(p.name());
}
@@ -467,7 +473,7 @@ class CallerResolverTest {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, members);
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
assertEquals(Role.OBSERVER, r.resolve("127.0.0.1", 42, null).role());
assertTrue(members.bind("architect:lead-designer", "term_a")); // the later lifecycle binds the slot
@@ -494,12 +500,13 @@ class CallerResolverTest {
}
@Test
void aBoundNonArchitectSlotStillResolvesAsAWorker() {
void aBoundNonArchitectSlotResolvesToTheObserverFloorNotArchitect() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, boundMembers("dev:builder", MemberRole.DEV))
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role(), "a dev binding must never grant architect rights");
assertEquals(Role.OBSERVER, p.role(), "a dev binding must never grant architect rights, "
+ "and this construction path wires no roster to recognise it as the live dev it is");
}
@Test
@@ -510,17 +517,17 @@ class CallerResolverTest {
assertThrows(IllegalArgumentException.class, () -> new CallerResolver(id, true, " "));
}
@Test
void aWorkerOnAnyLoopbackSourceAddressIsStillAWorkerNotThePrimary() {
// fleetd #305: the escalation. ConnectionIdentity used to accept only 127.0.0.1, so a
// worker connecting from 127.0.0.2 resolved to no terminal, and this resolver's own
// (wider) loopback check then made it the PRIMARY — granting spawn, stop, send and drain.
// Measured on the Linux fleet host: binding a source of 127.0.0.2 succeeds there, so the
// path is real and not theoretical.
void aPaneOnAnyLoopbackSourceAddressIsStillAtTheFloorNotThePrimary() {
// fleetd #305: the escalation this guards against. ConnectionIdentity used to accept only
// 127.0.0.1, so a pane connecting from 127.0.0.2 resolved to no terminal, and this
// resolver's own (wider) loopback check then made it the PRIMARY — granting spawn, stop,
// send and drain. Measured on the Linux fleet host: binding a source of 127.0.0.2 succeeds
// there, so the path is real and not theoretical.
CallerResolver r = new CallerResolver(workerIdentity(), false, null);
for (String src : new String[]{"127.0.0.1", "127.0.0.2", "127.1.2.3", "::ffff:127.0.0.2"}) {
Principal p = r.resolve(src, 55555, null);
assertEquals(Role.WORKER, p.role(), "a worker must stay a worker from source " + src);
assertEquals("term_a", p.terminal(), "worker terminal from source " + src);
assertEquals(Role.OBSERVER, p.role(), "the pane must stay off PRIMARY from source " + src);
assertEquals("term_a", p.terminal(), "pane terminal from source " + src);
}
}
@@ -826,14 +833,14 @@ class CallerResolverTest {
assertEquals("term_a", p.terminal());
}
/** Regression: an empty collaborator registry leaves every pane exactly as before. */
/** Regression: an empty collaborator registry leaves every pane at the unconfigured-pane floor. */
@Test
void anEmptyCollaboratorRegistryLeavesEveryPaneAsBefore() {
void anEmptyCollaboratorRegistryLeavesEveryPaneAtTheObserverFloor() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, Map::of)
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertNull(p.name());
}
@@ -862,6 +869,39 @@ class CallerResolverTest {
assertFalse(r.knownLeadOrCollaborator().test("term_other"));
}
// ── fleetd #705: narrowing the unconfigured-pane floor to OBSERVER ──────────────────────────
/**
* The case this ticket exists for: a pane the resolver cannot place as a live spawned member,
* a lead, a bound architect slot, or a configured collaborator must land on the narrow
* {@link Role#OBSERVER} floor, never the {@link Role#WORKER} the old fallback granted.
*
* <p>The second assertion is the control the ticket requires: a terminal the roster DOES
* recognise as a live spawned member must still resolve its own role. Without it, this test
* would also pass if the fix accidentally turned every caller into an observer.
*/
@Test
void anUnconfiguredPaneResolvesObserverButARegisteredMemberStillResolvesItsOwnRole() {
Principal unconfigured = new CallerResolver(workerIdentity()).resolve("127.0.0.1", 42, null);
assertEquals(Role.OBSERVER, unconfigured.role(),
"a pane matching none of the configured or live-roster roles must fall to the "
+ "floor, not WORKER");
assertEquals("term_a", unconfigured.terminal());
Principal registered = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, new MemberRegistry(null),
t -> "term_a".equals(t) ? MemberRole.DEV : null, Map::of)
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, registered.role(),
"control: a live spawned member must keep resolving its own role, never the "
+ "unconfigured-pane floor");
}
@Test
void describeNamesTheObserverByItsPane() {
assertEquals("observer:term_a", Principal.observer("term_a", 1).describe());
}
@Test
void knownLeadOrCollaboratorIsFalseForASpawnedMembersTerminal() {
// The exact scenario a collaborator's SEND must never reach: a live spawned member's own
@@ -295,9 +295,10 @@ class MemberRegistryLiveTest {
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
Principal after = resolver.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, after.role(),
"removing the slot from config must demote the bound session to worker on its "
+ "NEXT request — this is the ticket's whole point");
assertEquals(Role.OBSERVER, after.role(),
"removing the slot from config must demote the bound session on its NEXT request — "
+ "this harness wires no live roster for term_a, so the demotion lands on "
+ "the unconfigured-pane floor");
assertEquals("term_a", after.terminal(), "same pane, same terminal — only the role changed");
}
@@ -0,0 +1,31 @@
package dev.ltms.fleet.auth;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
class PrincipalTest {
@Test
void ownerKeyCoversEveryRole() {
assertEquals("leader:opus", Principal.leader("opus", "term_lead", 1).ownerKey());
assertNull(Principal.primary(2).ownerKey());
assertEquals("worker:term_worker", Principal.worker("term_worker", 3).ownerKey());
assertEquals("architect:term_arch", Principal.architect("opus", "term_arch", 4).ownerKey());
assertEquals("collaborator:ops", Principal.collaborator("ops", "term_collab", 5).ownerKey());
assertEquals("observer:term_observer", Principal.observer("term_observer", 6).ownerKey());
assertEquals("anonymous", Principal.anonymous().ownerKey());
}
@Test
void rolePrefixesKeepLeadAndArchitectKeysDistinct() {
String lead = Principal.leader("opus", "term_lead", 1).ownerKey();
String architect = Principal.architect("design", "opus", 2).ownerKey();
assertEquals("leader:opus", lead);
assertEquals("architect:opus", architect);
assertNotEquals(lead, architect);
}
}
@@ -103,7 +103,7 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
// comment there), same as broker/primary/leadHeartbeat/... above — a real, non-null value
// here proves it, rather than leaving it null and proving nothing.
v.put("leadRollover", new FleetConfig.LeadRollover(
"/handover/guard.md", true, 3600, 20, 20, "read the handover file"));
"/handover/guard.md", true, 3600, 20, 45, "read the handover file"));
assertNamesMatchComponents(v);
return v;
}
@@ -7,6 +7,7 @@ import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArrayList;
@@ -59,6 +60,10 @@ public final class FakeHerdr implements HerdrClient {
private Runnable onAgentStart; // fires the instant agent.start is called — see onAgentStart(Runnable)
private volatile int agentGetOkCalls = Integer.MAX_VALUE; // how many agent.get calls succeed first
private volatile String agentGetFailCode = null; // error code every agent.get call after that reports
/** pane ids that {@link #paneGoneAfterClose} has opted into reporting gone — see that method. */
private final Set<String> paneGoneAfterCloseIds = ConcurrentHashMap.newKeySet();
/** pane ids a {@code pane.close} call has actually reached, for {@link #paneGoneAfterCloseIds}. */
private final Set<String> closedPaneIds = ConcurrentHashMap.newKeySet();
public FakeHerdr healthy(boolean h) {
this.healthy = h;
@@ -220,6 +225,19 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/**
* Make {@code pane.get(paneId)} report the pane gone (a {@code pane_not_found} {@link
* HerdrException}, exactly as {@link WorkspaceControl#locatePane} expects to see once a pane
* has really disappeared) once a {@code pane.close} call for that same {@code paneId} has
* actually reached this fake. Every other pane, and this pane before its own close, keeps
* reporting the default canned {@code pane.get} response — opt-in, by pane id, so no existing
* test's {@code pane.get} behaviour changes.
*/
public FakeHerdr paneGoneAfterClose(String paneId) {
paneGoneAfterCloseIds.add(paneId);
return this;
}
/**
* Run {@code hook} synchronously the instant an {@code agent.start} call reaches this fake —
* i.e. the instant the peer PROCESS would start against a real herdr daemon. A test uses this
@@ -422,9 +440,18 @@ public final class FakeHerdr implements HerdrClient {
}
yield mapper.readTree("{\"type\":\"ok\"}");
}
case "pane.get" -> mapper.readTree("""
case "pane.get" -> {
Object paneIdParam = params instanceof Map<?, ?> m ? m.get("pane_id") : null;
String paneIdKey = paneIdParam == null ? null : String.valueOf(paneIdParam);
if (paneIdKey != null && paneGoneAfterCloseIds.contains(paneIdKey)
&& closedPaneIds.contains(paneIdKey)) {
throw new HerdrException("herdr error [pane_not_found]: pane.get failed",
"pane_not_found", null);
}
yield mapper.readTree("""
{"type":"pane_info","pane":{"pane_id":"w9:pW","workspace_id":"w9",
"tab_id":"w9:t2","agent_status":"idle"}}""");
}
case "pane.list" -> noPanes
? mapper.readTree("{\"type\":\"pane_list\",\"panes\":[]}")
: mapper.readTree("""
@@ -458,6 +485,9 @@ public final class FakeHerdr implements HerdrClient {
throw new HerdrException("herdr error [" + code + "]: pane.close failed",
code, null);
}
if (paneIdParam != null) {
closedPaneIds.add(String.valueOf(paneIdParam));
}
yield mapper.readTree("{\"type\":\"ok\"}");
}
default -> throw new HerdrException("fake has no canned response for " + method);
@@ -571,4 +571,137 @@ class LeadLauncherTest {
assertTrue(warn.contains("opus"), "names the profile: " + warn);
assertTrue(warn.contains("x".repeat(60)), "names the culprit argument: " + warn);
}
// ── fleetd #726 unit 1: the single-lead relaunch seam ─────────────────────────────────────
/**
* The returned agent's {@code terminalId()}/{@code paneId()} are the ones the fake
* {@code AgentControl} actually started — not a coincidental field left over from the caller.
* {@code paneId()} echoes the exact {@code pane_id} the launch's own {@code agent.start} call
* carried (protocol 19: the agent starts into the pane it is asked to), and {@code
* terminalId()} is herdr's own generated id, which the fake always shapes as {@code
* term_new_<n>}.
*/
@Test
void relaunchReturnsTheStartedAgent() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started, "a launchable, configured lead must start");
Object startedPaneIdParam = ((Map<?, ?>) herdr.lastCall("agent.start").params()).get("pane_id");
assertEquals(startedPaneIdParam, started.paneId(),
"paneId() must be the pane the agent.start call actually targeted");
assertTrue(started.terminalId() != null && started.terminalId().startsWith("term_new_"),
"terminalId() must be herdr's own generated id: " + started.terminalId());
}
/** The new tab is labelled with the lead's configured {@code tab:}, and AFTER the start. */
@Test
void relaunchLabelsTheNewTabAfterStarting() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started);
assertEquals("lead: opus", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"));
int startIndex = indexOfLastCall(herdr, "agent.start");
int renameIndex = indexOfLastCall(herdr, "tab.rename");
assertTrue(renameIndex > startIndex,
"the tab must be renamed AFTER the start succeeds, not before: start=" + startIndex
+ " rename=" + renameIndex);
}
private static int indexOfLastCall(FakeHerdr herdr, String method) {
int idx = -1;
List<FakeHerdr.Call> calls = herdr.calls;
for (int i = 0; i < calls.size(); i++) {
if (calls.get(i).method().equals(method)) {
idx = i;
}
}
return idx;
}
@Test
void relaunchOfAnUnknownLeadNameReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("not-declared");
assertNull(started);
assertFalse(herdr.called("agent.start"));
assertFalse(herdr.called("workspace.create"));
assertFalse(herdr.called("tab.create"));
}
@Test
void relaunchOfARecogniseOnlyLeadReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead(null, "lead: dead", 1))).relaunch("opus");
assertNull(started);
assertFalse(herdr.called("agent.start"));
}
@Test
void relaunchWithAnUnconfiguredProfileReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("nope", "lead: opus", 1))).relaunch("opus");
assertNull(started);
assertFalse(herdr.called("agent.start"));
}
/**
* The outer retry {@link LeadLauncher#relaunch(String)} owns, separate from {@code
* ResilientAgentLaunch}'s internal {@code agent_name_taken} retry: a failed attempt must not
* be the end of the whole relaunch. Each of the first two attempts exhausts {@code
* ResilientAgentLaunch.NAME_RETRIES} name attempts (every one of them rejected), so each
* attempt's own tab is created and then closed; the third attempt's first name is free.
*/
@Test
void relaunchRetriesTheWholeAttemptAndSucceedsOnTheThird() {
FakeHerdr herdr = new FakeHerdr()
.agentNameTakenTimes(2 * ResilientAgentLaunch.NAME_RETRIES);
dev.ltms.fleet.herdr.Agent started =
fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started, "the third attempt's first name is free — it must succeed");
assertEquals(3, herdr.calls.stream().filter(c -> c.method().equals("tab.create")).count(),
"one tab per attempt: three attempts");
assertEquals(2, herdr.calls.stream().filter(c -> c.method().equals("tab.close")).count(),
"the two failed attempts' tabs must be closed");
}
/**
* Every attempt fails outright (a herdr error {@code ResilientAgentLaunch} does not retry at
* all) — {@link LeadLauncher#relaunch(String)} must give up after exactly {@code
* RELAUNCH_ATTEMPTS} and must not leak any of the tabs it created along the way.
*/
@Test
void relaunchGivesUpAfterExactlyRelaunchAttemptsAndLeaksNoTab() {
FakeHerdr herdr = new FakeHerdr().agentStartFailsWith("some_other_error");
dev.ltms.fleet.herdr.Agent started =
fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNull(started, "every attempt failed — relaunch must give up, not hang or guess");
assertEquals(LeadLauncher.RELAUNCH_ATTEMPTS,
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"exactly RELAUNCH_ATTEMPTS attempts, no more, no fewer");
long tabsCreated = herdr.calls.stream().filter(c -> c.method().equals("tab.create")).count();
long tabsClosed = herdr.calls.stream().filter(c -> c.method().equals("tab.close")).count();
assertEquals(LeadLauncher.RELAUNCH_ATTEMPTS, tabsCreated);
assertEquals(tabsCreated, tabsClosed, "every tab this method created must be closed — no leaks");
}
}
File diff suppressed because it is too large Load Diff
@@ -504,12 +504,12 @@ class FleetMcpAuthzTest {
}
/**
* {@code fleet_poll{ticket}} must thread the calling connection's own terminal into
* {@code fleet_poll{ticket}} must thread the calling connection's owner key into
* {@link MessageService#poll(String, String)}, so a worker cannot read a ticket a different
* session created.
*/
@Test
void theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll() throws Exception {
void theFleetPollHandlerActuallyThreadsCallerOwnerIntoPoll() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("pollHandler =");
@@ -525,18 +525,18 @@ class FleetMcpAuthzTest {
"control failed: the scraped pollHandler block contains no poll(messages, ...) call "
+ "at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
"the fleet_poll handler must thread callerTerminal(exchange) into poll(...), not omit "
assertTrue(handlerBlock.contains("principal(exchange).ownerKey()"),
"the fleet_poll handler must thread principal(exchange).ownerKey() into poll(...), not omit "
+ "it or pass a literal null -- block: " + handlerBlock);
}
/**
* {@code fleet_status} must thread the calling connection's own terminal into
* {@code fleet_status} must thread the calling connection's owner key into
* {@link FleetMcp#status(MessageService, String, String)}, so a caller that did not create a
* worker's open delegation cannot read its pending question through the status handler either.
*/
@Test
void theFleetStatusHandlerActuallyThreadsCallerTerminalIntoStatus() throws Exception {
void theFleetStatusHandlerActuallyThreadsCallerOwnerIntoStatus() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("statusHandler =");
@@ -552,8 +552,8 @@ class FleetMcpAuthzTest {
"control failed: the scraped statusHandler block contains no status(messages, ...) "
+ "call at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
"the fleet_status handler must thread callerTerminal(exchange) into status(...), not "
assertTrue(handlerBlock.contains("principal(exchange).ownerKey()"),
"the fleet_status handler must thread principal(exchange).ownerKey() into status(...), not "
+ "omit it or pass a literal null -- block: " + handlerBlock);
}
@@ -11,6 +11,7 @@ import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.lead.LeadRollover;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
@@ -24,6 +25,8 @@ import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
@@ -36,14 +39,14 @@ import static org.junit.jupiter.api.Assertions.*;
* fleetd #480 Unit C — the {@code fleet_handover} MCP tool, the surface that finally calls
* {@link LeadRollover#open}/{@link LeadRollover#confirm}/{@link LeadRollover#cancel}.
*
* <p>Uses {@link LeadRollover}'s PUBLIC constructor (real wall clock, real 250ms settle poll, a
* real virtual-thread continuation runner) rather than its package-private test constructor —
* this test lives in {@code dev.ltms.fleet.mcp}, not {@code dev.ltms.fleet.lead}, and does not
* need to control the post-{@code confirm()} continuation's timing: it only asserts the
* SYNCHRONOUS return value of {@code open}/{@code confirm}/{@code cancel}, which is exactly what
* {@code FleetMcp.handover} forwards to the client. {@code turnSettleSeconds}/{@code
* clearSettleSeconds} are kept at 1s so a confirmed request's background continuation (which this
* class does not wait on or assert against) gives up quickly rather than polling for 20s on a
* <p>Uses {@link LeadRollover}'s PUBLIC constructor (real wall clock, real 250ms poll, a real
* virtual-thread continuation runner) rather than its package-private test constructor — this
* test lives in {@code dev.ltms.fleet.mcp}, not {@code dev.ltms.fleet.lead}, and does not need to
* control the post-{@code confirm()} continuation's timing: it only asserts the SYNCHRONOUS return
* value of {@code open}/{@code confirm}/{@code cancel}, which is exactly what {@code
* FleetMcp.handover} forwards to the client. {@code turnSettleSeconds}/{@code
* relaunchReadySeconds} are kept at 1s so a confirmed request's background continuation (which
* this class does not wait on or assert against) gives up quickly rather than polling for 20s on a
* daemon virtual thread.
*/
class FleetMcpHandoverTest {
@@ -67,10 +70,26 @@ class FleetMcpHandoverTest {
return new FleetConfig.LeadRollover(handoverPath, false, 3600, 1, 1, "read the handover file");
}
private static FleetConfig minimalFleetConfig() {
try {
Path yaml = Files.createTempFile("fleet-mcp-handover-test", ".yaml");
Files.writeString(yaml, "bind:\n port: 8080\n");
return FleetConfig.load(yaml);
} catch (IOException e) {
throw new UncheckedIOException(e);
}
}
private LeadRollover newRollover(String handoverPath) {
// Every handoverPath this test class uses comes from tmp.resolve(...), which is already
// absolute, so the workspace lookup is never actually consulted — a no-op lookup is enough.
return new LeadRollover(agents, () -> cfg(handoverPath), _ -> null);
// None of this class's tests reach the recognition-wait or the relaunch call, so the
// launcher's own functional correctness is irrelevant here — any constructed instance,
// backed by the same fake herdr, is enough.
WorkspaceControl spaces = new WorkspaceControl(herdr);
LeadLauncher launcher = new LeadLauncher(agents, spaces, minimalFleetConfig());
return new LeadRollover(agents, spaces, launcher, () -> cfg(handoverPath),
_ -> null, _ -> null, Map::of);
}
/** A fully wired FleetMcp on fakes (mirrors FleetMcpAuthzTest's helper), plus a leadRollover. */
@@ -329,12 +329,13 @@ class FleetMcpTest {
/**
* Same hijack and control through the fire-and-poll ({@code sendAsync}) path: the owner comes
* from the ticket's recorded creator terminal, not from a caller threaded through a live call.
* from the resolved principal, not from a caller argument threaded through a live call.
*/
@Test
void aDifferentCallersMcpAnswerIsRefusedForAnAsyncSendButTheRealOwnerSucceeds() throws Exception {
Principal owner = Principal.worker("term_owner", 1);
McpSchema.CallToolResult accepted =
FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), "term_owner");
FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), owner);
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
@@ -354,14 +355,15 @@ class FleetMcpTest {
assertEquals(MessageService.Phase.ASKING, asking.phase());
String turnId = asking.turnId();
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L,
"worker:term_attacker");
assertTrue(hijacked.isError(), "a caller that did not create this delegation must get an error");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
"a refused answer must not advance the async ticket's phase");
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey()));
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -1920,8 +1922,8 @@ class FleetMcpTest {
Principal architect = Principal.architect("lead-designer", "term_design", 400);
MemberPresence presence = new MemberPresence();
FleetMcp.markSpawnedMemberPresent(worker, presence);
FleetMcp.markSpawnedMemberPresent(architect, presence);
FleetMcp.markTrackedCallerPresent(worker, presence);
FleetMcp.markTrackedCallerPresent(architect, presence);
assertTrue(presence.isPresent("term_worker"));
assertTrue(presence.isPresent("term_design"));
@@ -1932,12 +1934,27 @@ class FleetMcpTest {
Principal lead = Principal.leader("opus", "term_lead", 100);
MemberPresence presence = new MemberPresence();
FleetMcp.markSpawnedMemberPresent(lead, presence);
FleetMcp.markSpawnedMemberPresent(Principal.anonymous(), presence);
FleetMcp.markTrackedCallerPresent(lead, presence);
FleetMcp.markTrackedCallerPresent(Principal.anonymous(), presence);
assertFalse(presence.isPresent("term_lead"));
}
/**
* The item whose absence would be silent: an observer's own MCP contact must still mark
* presence, or a pane resolving to the unconfigured-pane floor would sit on the injector
* readiness gate forever once something addresses it.
*/
@Test
void anObserverContactMarksPresence() {
Principal observer = Principal.observer("term_observer", 800);
MemberPresence presence = new MemberPresence();
FleetMcp.markTrackedCallerPresent(observer, presence);
assertTrue(presence.isPresent("term_observer"));
}
@Test
void statusReportsLiveAgentStatus() {
FakeHerdr blocked = new FakeHerdr().agentStatus("blocked");
@@ -1996,13 +2013,13 @@ class FleetMcpTest {
/**
* {@code fleet_status}'s pending-ask block (the question, its {@code turnId} and its ticket)
* is shown only to the caller whose terminal created the delegation, or to a caller with no
* terminal at all (the unnamed primary) — a different terminal-bearing caller still sees the
* base status line, but none of the pending-ask fields.
* is shown only to the caller whose owner key created the delegation, or to the unnamed primary.
* A different caller still sees the base status line, but none of the pending-ask fields.
*/
@Test
void statusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, "term_creator");
void statusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
Principal creator = Principal.worker("term_creator", 1);
String ticket = messages.sendAsync(T, "task that asks", null, creator);
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -2020,7 +2037,7 @@ class FleetMcpTest {
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
String other = textOf(FleetMcp.status(messages, T, "term_other"));
String other = textOf(FleetMcp.status(messages, T, "worker:term_other"));
assertTrue(other.startsWith("idle"), "the base status must still be shown: " + other);
assertFalse(other.contains("which config file?"),
"a non-creating caller must not see the question text: " + other);
@@ -2029,10 +2046,12 @@ class FleetMcpTest {
assertFalse(other.contains(ticket),
"a non-creating caller must not see the ticket: " + other);
String creator = textOf(FleetMcp.status(messages, T, "term_creator"));
assertTrue(creator.contains("which config file?"), "the creator must see the question: " + creator);
assertTrue(creator.contains(asking.turnId()), "the creator must see the turnId: " + creator);
assertTrue(creator.contains(ticket), "the creator must see the ticket: " + creator);
String creatorStatus = textOf(FleetMcp.status(messages, T, creator.ownerKey()));
assertTrue(creatorStatus.contains("which config file?"),
"the creator must see the question: " + creatorStatus);
assertTrue(creatorStatus.contains(asking.turnId()),
"the creator must see the turnId: " + creatorStatus);
assertTrue(creatorStatus.contains(ticket), "the creator must see the ticket: " + creatorStatus);
String unnamed = textOf(FleetMcp.status(messages, T, null));
assertTrue(unnamed.contains("which config file?"),
@@ -2041,7 +2060,7 @@ class FleetMcpTest {
// Clean up the still-open ask so the background thread does not linger past the test.
String turnId = asking.turnId();
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(turnId, "config.yaml", 5000, "term_creator"));
() -> messages.answer(turnId, "config.yaml", 5000, creator.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -2149,6 +2168,27 @@ class FleetMcpTest {
assertTrue(leadOut.contains("\"leader\":\"opus\""), leadOut);
}
/**
* An observer reports its own role and pane, never a {@code leader} key. Without an explicit
* branch it would reach the lead branch by elimination and look right only because the
* {@code leader} key is guarded on a non-null name — this pins the branch rather than the
* accident.
*/
@Test
void whoamiReportsAnObserverNotALead() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw"));
McpSchema.CallToolResult res = FleetMcp.whoami(
Principal.observer("term_observer", 900), sessions);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"role\":\"observer\""), out);
assertTrue(out.contains("\"sessionId\":\"term_observer\""), out);
assertFalse(out.contains("leader"), out);
}
/**
* CB-548: an architect SEND delegates as its own pane (recording the per-target delegation) but
* must NEVER become the legacy singleton "primary" fallback — the per-target map does not cure
@@ -19,7 +19,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
* {@link MessageService#poll(String)} overload. That overload skips the ownership check in
* {@code MessageService}'s {@code ownsTicket} entirely, so a caller of it can read any session's
* ticket. Every production caller must go through {@link MessageService#poll(String, String)}
* and pass a {@code callerTerminal} explicitly, even when it is {@code null}.
* and pass a {@code callerOwner} explicitly, even when it is {@code null}.
*
* <p>This reads each file's own source text rather than reflecting on compiled bytecode, because
* the risk is a future one-word edit at a call site, not a missing overload.
@@ -48,7 +48,7 @@ class MessageServicePollUsageTest {
+ "below proves nothing");
assertTrue(violations.isEmpty(), "found a call to the fail-open MessageService.poll(String) "
+ "overload, which skips the ownership check entirely -- pass a callerTerminal "
+ "overload, which skips the ownership check entirely -- pass a callerOwner "
+ "explicitly (even if null) through poll(String, String) instead: " + violations);
// CONTROL: the arity parser actually finds the two genuine two-argument call sites (the
@@ -58,7 +58,7 @@ class MessageServicePollUsageTest {
assertEquals(2, twoArgSites.size(), "control failed: expected exactly the two known "
+ "two-argument messages.poll(...) call sites, found: " + twoArgSites);
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetMcp.java")),
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerTerminal) "
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerOwner) "
+ "site among: " + twoArgSites);
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetApp.java")),
"control failed: did not find the FleetApp.java messages.poll(...) site among: "
@@ -4,6 +4,7 @@ import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
@@ -888,24 +889,59 @@ class MessageServiceTest {
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
}
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
// --- fleetd #719: a per-boot nonce keeps one instance's ticket ids out of another's space ---
/** A second, fully independent instance — its own agents/injector/rendezvous/inbox, not shared. */
private MessageService newIndependentInstance() {
FakeHerdr otherHerdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
AgentControl otherAgents = new AgentControl(otherHerdr);
Injector otherInjector = new Injector(otherAgents);
return new MessageService(otherAgents, otherInjector, new Rendezvous(), new InMemoryReplyInbox());
}
@Test
void pollByAnotherTerminalIsRefused() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_a");
void twoInstancesMintDisjointTicketIds() {
MessageService other = newIndependentInstance();
String ticketFromThis = messages.sendAsync(T, "task on first instance", null, null);
String ticketFromOther = other.sendAsync(T, "task on second instance", null, null);
assertNotEquals(ticketFromThis, ticketFromOther,
"each instance mints its own id space, so even a first ticket from each must differ");
}
@Test
void foreignInstanceTicketDoesNotResolve() {
MessageService other = newIndependentInstance();
String ticket = messages.sendAsync(T, "task on first instance", null, null);
// `other` must reach the same sequence number, or this test passes against an empty map
// instead of against a colliding id.
other.sendAsync(T, "task on second instance", null, null);
// control: the id resolves in the instance that minted it, so a null below cannot be
// explained by broken plumbing — only by the ticket being foreign to `other`.
assertNotNull(messages.poll(ticket), "the minting instance must still resolve its own ticket");
assertNull(other.poll(ticket), "a ticket minted by a different instance must not resolve here");
}
// --- ticket ownership -----------------------------------------------------------------------
@Test
void leadBIsRefusedFromLeadAsTicket() throws Exception {
Principal leadA = Principal.leader("opus", "term_a", 1);
Principal leadB = Principal.leader("sol", "term_b", 2);
String ticket = messages.sendAsync(T, "long task", null, leadA);
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "secret async result"), "a reply resolves the async send");
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, "term_a");
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, leadA.ownerKey());
assertNotNull(owner, "the creator must still be able to read its own ticket");
assertEquals(MessageService.Phase.DONE, owner.phase());
MessageService.TaskView refused = messages.poll(ticket, "term_b");
assertNotNull(refused, "a different terminal gets a refusal, not silence");
MessageService.TaskView refused = messages.poll(ticket, leadB.ownerKey());
assertNotNull(refused, "a different lead gets a refusal, not silence");
assertNotEquals(MessageService.Phase.DONE, refused.phase(),
"a different terminal must never see the ticket as DONE");
"a different lead must never see the ticket as DONE");
assertNull(refused.reply(), "a refusal must never carry the reply text");
assertFalse(String.valueOf(refused).contains("secret async result"),
"the reply text must not appear anywhere in the refused view");
@@ -913,14 +949,13 @@ class MessageServiceTest {
@Test
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
String ticket = messages.sendAsync(T, "long task", null,
Principal.leader("opus", "term_lead", 1));
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
// callerTerminal == null is the unnamed primary (resolved by token or loopback trust, with
// no herdr pane) — it must read a ticket a terminal-bearing lead created.
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
assertNotNull(view, "the unnamed primary must be able to read any ticket");
assertEquals(MessageService.Phase.DONE, view.phase());
@@ -928,19 +963,44 @@ class MessageServiceTest {
}
@Test
void creatorReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
void namedLeadCanPollItsTicketAfterItsTerminalChanges() throws Exception {
Principal oldLead = Principal.leader("opus", "term_OLD", 1);
Principal newLead = Principal.leader("opus", "term_NEW", 2);
assertNotEquals(oldLead.terminal(), newLead.terminal(), "the test requires different terminals");
String ticket = messages.sendAsync(T, "long task", null, oldLead);
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "own result"), "a reply resolves the async send");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, "term_creator");
assertNotNull(view, "the ticket's own creator must be able to read it");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, newLead.ownerKey());
assertNotNull(view, "the same named lead must read the ticket from its new terminal");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("own result", view.reply());
}
@Test
void anonymousOwnerKeyIsRefusedByTheTicketGateItself() {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "long task", null, lead);
MessageService.TaskView refused = messages.poll(ticket, Principal.anonymous().ownerKey());
assertNotNull(refused);
assertEquals(MessageService.Phase.FAILED, refused.phase());
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
}
@Test
void architectOwnershipUsesTerminalRatherThanSlot() {
Principal oldArchitect = Principal.architect("opus", "term_OLD", 1);
Principal newArchitect = Principal.architect("opus", "term_NEW", 2);
String ticket = messages.sendAsync(T, "long task", null, oldArchitect);
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket, oldArchitect.ownerKey()).phase());
assertEquals(MessageService.Phase.FAILED, messages.poll(ticket, newArchitect.ownerKey()).phase());
}
@Test
void pollReportsACompletedTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task");
@@ -1917,13 +1977,13 @@ class MessageServiceTest {
}
/**
* A caller's own terminal must match the terminal that created the delegation to see its
* pending question; a different terminal-bearing caller sees nothing, and a caller with no
* terminal at all (the unnamed primary) always sees it.
* A caller's owner key must match the key that created the delegation to see its pending
* question. The unnamed primary always sees it.
*/
@Test
void pendingAskGatesTheQuestionByTheDelegationsCreatorTerminal() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, "term_creator");
void pendingAskGatesTheQuestionByTheDelegationsCreatorOwner() throws Exception {
Principal creator = Principal.worker("term_creator", 1);
String ticket = messages.sendAsync(T, "task that asks", null, creator);
awaitWaiting();
injectDelivery();
@@ -1931,10 +1991,10 @@ class MessageServiceTest {
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertNull(messages.pendingAsk(T, "term_other"),
"a caller whose terminal did not create the delegation must not see the question");
assertNull(messages.pendingAsk(T, "worker:term_other"),
"a caller whose key did not create the delegation must not see the question");
MessageService.PendingAsk own = messages.pendingAsk(T, "term_creator");
MessageService.PendingAsk own = messages.pendingAsk(T, creator.ownerKey());
assertNotNull(own, "the creating caller must see its own open question");
assertEquals("which config file?", own.question());
@@ -1943,7 +2003,7 @@ class MessageServiceTest {
assertEquals("which config file?", unnamed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_creator"));
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1951,13 +2011,12 @@ class MessageServiceTest {
}
/**
* A task created with no recorded creator terminal (a short {@code sendAsync} overload) must
* not hand its open question to any caller that does have a terminal — only a caller with no
* terminal at all may still see it.
* A task created with no recorded owner (a short {@code sendAsync} overload) must not hand its
* open question to a caller with an owner key. Only the unnamed primary may still see it.
*/
@Test
void pendingAskDeniesATerminalBearingCallerWhenTheTaskRecordsNoCreator() throws Exception {
String ticket = messages.sendAsync(T, "task that asks"); // no creatorTerminal recorded
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
@@ -1965,8 +2024,8 @@ class MessageServiceTest {
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertNull(messages.pendingAsk(T, "term_someone"),
"a terminal-bearing caller must not see a question whose task records no creator");
assertNull(messages.pendingAsk(T, "worker:term_someone"),
"a caller with an owner key must not see a question whose task records no creator");
assertNotNull(messages.pendingAsk(T, null),
"the unnamed primary must still see it even with no recorded creator");
@@ -2022,7 +2081,8 @@ class MessageServiceTest {
*/
@Test
void aDifferentCallersAnswerIsRefusedForAnAsyncSendDelegationButTheRealOwnerSucceeds() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, "term_owner");
Principal owner = Principal.worker("term_owner", 1);
String ticket = messages.sendAsync(T, "task that asks", null, owner);
awaitWaiting();
injectDelivery();
@@ -2030,7 +2090,8 @@ class MessageServiceTest {
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Reply hijacked = messages.answer(asking.turnId(), "evil.yaml", 500, "term_attacker");
MessageService.Reply hijacked = messages.answer(asking.turnId(), "evil.yaml", 500,
"worker:term_attacker");
assertEquals(MessageService.Outcome.NOT_TURN_OWNER, hijacked.outcome(),
"a caller that did not create this delegation must be refused, not served");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
@@ -2038,7 +2099,38 @@ class MessageServiceTest {
"a refused answer must not advance the async ticket's phase");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_owner"));
() -> messages.answer(asking.turnId(), "config.yaml", 5000, owner.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
assertEquals("done", awaitTicketPhase(ticket, MessageService.Phase.DONE).reply());
}
@Test
void namedLeadCanSeeAndAnswerAnAskAfterItsTerminalChangesWhileLeadBIsRefused() throws Exception {
Principal oldLead = Principal.leader("opus", "term_OLD", 1);
Principal newLead = Principal.leader("opus", "term_NEW", 2);
Principal leadB = Principal.leader("sol", "term_SOL", 3);
assertNotEquals(oldLead.terminal(), newLead.terminal(), "the test requires different terminals");
String ticket = messages.sendAsync(T, "task that asks", null, oldLead);
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertNotNull(messages.pendingAsk(T, newLead.ownerKey()),
"the same lead at its new terminal must see the pending ask");
assertNull(messages.pendingAsk(T, leadB.ownerKey()),
"another lead must not see the pending ask");
assertEquals(MessageService.Outcome.NOT_TURN_OWNER,
messages.answer(asking.turnId(), "evil.yaml", 500, leadB.ownerKey()).outcome());
assertFalse(ask.isDone(), "another lead must not resolve the worker's ask");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, newLead.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -191,21 +191,68 @@ class RendezvousTest {
assertNull(rendezvous.askOwner(t.turnId()), "a closed ask no longer reports an owner");
}
// ── fleetd #729: per-boot nonce guards turnId against cross-instance reuse ────────────────
@Test
void twoInstancesMintDisjointTurnIds() {
Rendezvous other = new Rendezvous();
Rendezvous.AskTicket fromThis = rendezvous.openAsk(W);
Rendezvous.AskTicket fromOther = other.openAsk(W);
assertNotEquals(fromThis.turnId(), fromOther.turnId(),
"each instance mints its own id space, so even a first ask from each must differ");
}
@Test
void foreignInstanceTurnIdDoesNotResolve() {
Rendezvous other = new Rendezvous();
// `other` must reach the same sequence number as `rendezvous` (two asks each, the first
// closed so the second mints fresh), or this test passes against an empty map instead of
// against a colliding id.
Rendezvous.AskTicket firstFromThis = rendezvous.openAsk(W);
rendezvous.closeAsk(firstFromThis.turnId());
Rendezvous.AskTicket secondFromThis = rendezvous.openAsk(W);
Rendezvous.AskTicket firstFromOther = other.openAsk(W);
other.closeAsk(firstFromOther.turnId());
other.openAsk(W);
// control: the id resolves in the instance that minted it, so a false below cannot be
// explained by broken plumbing — only by the turnId being foreign to `other`.
assertTrue(rendezvous.answerAsk(secondFromThis.turnId(), "answer from this instance"),
"the minting instance must still resolve its own turnId");
assertFalse(other.answerAsk(secondFromThis.turnId(), "answer from other instance"),
"a turnId minted by a different instance must not resolve here");
}
@Test
void openAskStillCoalescesDuplicatesAndStillMintsDistinctIdsPerAsk() {
Rendezvous.AskTicket t1 = rendezvous.openAsk(W);
Rendezvous.AskTicket t2 = rendezvous.openAsk(W);
assertEquals(t1.turnId(), t2.turnId(),
"a second openAsk while one is open still coalesces onto the same turn");
assertFalse(t2.fresh(), "the coalesced ask is still reported as not fresh");
rendezvous.closeAsk(t1.turnId());
Rendezvous.AskTicket t3 = rendezvous.openAsk(W);
assertNotEquals(t1.turnId(), t3.turnId(), "two asks from the same session still get different turnIds");
}
@Test
void ownerPermitsIsFailClosedOnARecordAndThreeStatesAreDistinct() {
assertFalse(Rendezvous.Owner.permits(null, null),
"no owner on record refuses even a caller with no terminal");
assertFalse(Rendezvous.Owner.permits(null, "term_a"),
"no owner on record refuses a terminal-bearing caller too");
"no owner on record refuses even the unnamed primary");
assertFalse(Rendezvous.Owner.permits(null, "worker:term_a"),
"no owner on record refuses a caller with an owner key too");
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, null),
"the unnamed primary owner matches a caller with no terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, "term_a"),
"the unnamed primary owner does not match a terminal-bearing caller");
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_a"),
"a named owner matches the same terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_b"),
"a named owner refuses a different terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), null),
"the recorded unnamed primary matches a caller with a null owner key");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, "worker:term_a"),
"the recorded unnamed primary does not match another owner key");
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), "worker:term_a"),
"an owner matches the same key");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), "worker:term_b"),
"an owner refuses a different key");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), null),
"a named owner refuses the unnamed primary");
}
}
@@ -8,6 +8,7 @@ import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.Authz;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.auth.Role;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
@@ -20,6 +21,7 @@ import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.FakeWorktrees;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
@@ -146,11 +148,11 @@ class FleetAppAuthTest {
/**
* {@code GET /tasks/{ticket}} must resolve its caller the same way {@code allow(...)} does
* and thread that terminal into {@link MessageService#poll(String, String)}, not the
* and thread that owner key into {@link MessageService#poll(String, String)}, not the
* no-check overload that ignores who is asking.
*/
@Test
void theTaskStatusRouteActuallyThreadsTheCallersTerminalIntoPoll() throws Exception {
void theTaskStatusRouteActuallyThreadsTheCallersOwnerKeyIntoPoll() throws Exception {
String source = Files.readString(REST_SOURCE);
int start = source.indexOf("private void taskStatus(Context ctx) {");
@@ -166,8 +168,8 @@ class FleetAppAuthTest {
"control failed: the scraped taskStatus block contains no messages.poll( call at all "
+ "-- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("caller.terminal()"),
"the taskStatus route must thread the resolved caller's terminal into messages.poll(...), "
assertTrue(handlerBlock.contains("caller.ownerKey()"),
"the taskStatus route must thread the resolved caller's owner key into messages.poll(...), "
+ "not the no-check overload -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
"the taskStatus route must resolve its caller the same way allow(...) does, not via a "
@@ -176,11 +178,11 @@ class FleetAppAuthTest {
/**
* {@code GET /sessions/{id}/status} must resolve its caller the same way {@code allow(...)}
* does and thread that terminal into {@link MessageService#pendingAsk(String, String)}, not
* does and thread that owner key into {@link MessageService#pendingAsk(String, String)}, not
* the no-check overload that ignores who is asking.
*/
@Test
void theSessionStatusRouteActuallyThreadsTheCallersTerminalIntoPendingAsk() throws Exception {
void theSessionStatusRouteActuallyThreadsTheCallersOwnerKeyIntoPendingAsk() throws Exception {
String source = Files.readString(REST_SOURCE);
int start = source.indexOf("private void sessionStatus(Context ctx) {");
@@ -196,8 +198,8 @@ class FleetAppAuthTest {
"control failed: the scraped sessionStatus block contains no messages.pendingAsk( "
+ "call at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("caller.terminal()"),
"the sessionStatus route must thread the resolved caller's terminal into "
assertTrue(handlerBlock.contains("caller.ownerKey()"),
"the sessionStatus route must thread the resolved caller's owner key into "
+ "messages.pendingAsk(...), not the no-check overload -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
"the sessionStatus route must resolve its caller the same way allow(...) does, not via "
@@ -222,7 +224,8 @@ class FleetAppAuthTest {
Javalin otherWorkerApp = startOnSharedService(messages, herdr, 9001L); // -> term_shell
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
try {
String ticket = messages.sendAsync("term_a", "long task", null, "term_a");
String ticket = messages.sendAsync("term_a", "long task", null,
Principal.worker("term_a", FakeHerdr.WORKER_PID));
HttpResponse<String> refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, refused.statusCode());
@@ -287,12 +290,12 @@ class FleetAppAuthTest {
/**
* {@code GET /sessions/{id}/status} shows a worker's pending {@code fleet_ask} question, its
* {@code turnId} and its ticket only to the caller whose terminal created that delegation, or
* to a caller with no terminal at all (the unnamed primary) — a different terminal-bearing
* caller still sees the base status line, but none of the pending-ask fields.
* {@code turnId} and its ticket only to the caller whose owner key created that delegation, or
* to the unnamed primary. A different caller still sees the base status line, but none of the
* pending-ask fields.
*/
@Test
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception {
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
@@ -304,7 +307,8 @@ class FleetAppAuthTest {
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
try {
ObjectMapper mapper = new ObjectMapper();
String ticket = messages.sendAsync("term_target", "task that asks", null, "term_a");
String ticket = messages.sendAsync("term_target", "task that asks", null,
Principal.worker("term_a", FakeHerdr.WORKER_PID));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -343,7 +347,7 @@ class FleetAppAuthTest {
// Clean up the still-open ask so the background thread does not linger past the test.
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(turnId, "config.yaml", 5000, "term_a"));
() -> messages.answer(turnId, "config.yaml", 5000, "worker:term_a"));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
@@ -520,6 +524,11 @@ class FleetAppAuthTest {
* As {@link #startOnSharedService(MessageService, FakeHerdr, long)}, but {@code leadTerminals}
* resolves the given pid's terminal to a named lead (a caller with SEND permission) instead of
* a plain worker, for a test that needs a terminal-bearing caller able to create a ticket.
*
* <p>Every connecting pane not already claimed by {@code leadTerminals} is wired into the live
* roster as a spawned worker, so a caller's resolved role matches what its own test expects:
* a {@link Role#WORKER}, never the unconfigured-pane {@link Role#OBSERVER} floor a roster-less
* resolver would otherwise fall to.
*/
private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid,
Map<String, String> leadTerminals) {
@@ -534,7 +543,8 @@ class FleetAppAuthTest {
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> leadTerminals, new MemberRegistry(null));
() -> leadTerminals, new MemberRegistry(null),
t -> leadTerminals.containsKey(t) ? null : MemberRole.DEV, Map::of);
Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
@@ -0,0 +1,92 @@
package dev.ltms.fleet.session;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.PlacementDecision;
import java.util.List;
import java.util.Set;
/**
* {@link PeerLauncher} decorator that marks presence for a spawned terminal before returning its
* handle to the caller — the contact-then-register ordering fleetd #722 covers, where the
* terminal's MCP contact lands before {@link SessionManager#acquire} runs its own
* {@code registry.put}. The presence view is set after construction, via {@link #presence},
* because it is owned by the {@link SessionManager} this launcher is passed into.
*/
final class PresenceRacingLauncher implements PeerLauncher {
private final PeerLauncher delegate;
volatile MemberPresence presence;
PresenceRacingLauncher(PeerLauncher delegate) {
this.delegate = delegate;
}
@Override
public PeerHandle spawn(SpawnRequest req) {
PeerHandle handle = delegate.spawn(req);
presence.markPresent(handle.terminalId());
return handle;
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
PeerHandle handle = delegate.spawn(req, decision);
presence.markPresent(handle.terminalId());
return handle;
}
@Override
public Set<Capability> capabilities() {
return delegate.capabilities();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return delegate.capabilitiesFor(profileName);
}
@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) {
delegate.stop(id);
}
@Override
public boolean clearContext(String id) {
return delegate.clearContext(id);
}
}
@@ -422,6 +422,53 @@ class SessionManagerTest {
"turn completion moves BUSY → DONE");
}
// --- fleetd #722: registration and presence must reach READY whichever lands first --------
@Test
void registerThenContactReachesReadyForPlainSpawn() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
"a presence contact that arrives after registration reaches READY");
}
@Test
void contactThenRegisterStillReachesReadyForPlainSpawn() {
// The racing launcher marks presence for the spawned terminal from inside spawn() —
// before SessionManager.acquire's own registry.put runs — modeling an MCP contact that
// lands in that window.
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 workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
PresenceRacingLauncher race = new PresenceRacingLauncher(workers);
SessionManager sessions = new SessionManager(race);
race.presence = sessions.asPresence();
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
"a presence contact that lands before registry.put must still reach READY");
}
@Test
void aTerminalNeverMarkedPresentStaysSpawningAfterRegistration() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
assertEquals(MemberSession.State.SPAWNING, sessions.get(session.paneId()).orElseThrow().state(),
"registration alone must not advance a terminal that was never marked present");
}
@Test
void releaseTearsDownWorkerAndRemovesFromRosterAndIsIdempotent() {
FakeHerdr herdr = new FakeHerdr();
@@ -1405,6 +1452,105 @@ class SessionManagerTest {
+ "dirty check threw");
}
// --- fleetd #736: a release must forget the member's presence entry, not just its registry
// row ---------------------------------------------------------------------------------------
@Test
void releaseByPaneIdForgetsThePresenceEntry() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
assertTrue(sessions.asPresence().isPresent(terminal), "present before the release");
sessions.release(session.paneId());
assertFalse(sessions.asPresence().isPresent(terminal),
"release must forget the terminal's presence, not just remove its registry row");
}
@Test
void reapIdleForgetsThePresenceEntryToo() {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
assertTrue(sessions.asPresence().isPresent(terminal), "present before the reap");
clock[0] = 11;
assertEquals(1, sessions.reapIdle(10), "READY session past TTL is reaped");
assertFalse(sessions.asPresence().isPresent(terminal),
"the idle-reap release path (releaseIfCurrent) goes through the same teardown "
+ "funnel as an explicit release, so it must forget presence too");
}
@Test
void shutdownDrainAlsoForgetsThePresenceEntry() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
assertTrue(sessions.asPresence().isPresent(terminal), "present before the drain");
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
assertFalse(sessions.asPresence().isPresent(terminal),
"a shutdown drain still ends the member's process, so presence must be cleared "
+ "exactly as it is for any other release cause");
}
@Test
void releaseOfAnUnknownPaneIdDoesNotThrow() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
assertDoesNotThrow(() -> sessions.release("no-such-pane"),
"releasing a pane id that was never registered must be a no-op, not a throw");
}
@Test
void releaseStillForgetsPresenceWhenDirtyCheckThrows() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("fleetd-736", null));
String terminal = s.terminalId();
sessions.asPresence().markPresent(terminal);
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
assertDoesNotThrow(() -> sessions.release(s.paneId()),
"a throwing dirty check must not abort the release");
assertFalse(sessions.asPresence().isPresent(terminal),
"presence must be forgotten even when the dirty check throws, which pins the "
+ "forget call to the finally block that runs no matter what happened above");
}
@Test
void releaseLeavesADifferentStillLiveMembersPresenceUntouched() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession released = sessions.acquire("ltms-local", null, "/caller/a", "ownerA");
MemberSession stillLive = sessions.acquire("ltms-local", null, "/caller/b", "ownerB");
sessions.asPresence().markPresent(released.terminalId());
sessions.asPresence().markPresent(stillLive.terminalId());
assertTrue(sessions.asPresence().isPresent(stillLive.terminalId()),
"present before the release of the other member");
sessions.release(released.paneId());
assertFalse(sessions.asPresence().isPresent(released.terminalId()),
"the released terminal is forgotten");
assertTrue(sessions.asPresence().isPresent(stillLive.terminalId()),
"a still-live member's presence must survive an unrelated release");
}
// --- fleetd #316: the dirty check must be re-taken after the worker is stopped, not trusted
// stale from before it ------------------------------------------------------------------------
@@ -159,6 +159,40 @@ class WorktreeSessionManagerTest {
assertEquals(expectedPath, s.cwd(), "session cwd is the worktree path");
}
// --- fleetd #722: registration and presence must reach READY whichever lands first --------
@Test
void registerThenContactReachesReadyForWorktreeSpawn() {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
MemberSession session = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
new WorktreeRequest("cb-722", null));
sessions.asPresence().markPresent(session.terminalId());
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
"a presence contact that arrives after worktree registration reaches READY");
}
@Test
void contactThenRegisterStillReachesReadyForWorktreeSpawn() {
// The racing launcher marks presence for the spawned terminal from inside spawn() —
// before SessionManager.acquireWithWorktree's own registry.put runs — modeling an MCP
// contact that lands in that window.
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
PresenceRacingLauncher race = new PresenceRacingLauncher(workerService(herdr));
SessionManager sessions = new SessionManager(race, worktrees);
race.presence = sessions.asPresence();
MemberSession session = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
new WorktreeRequest("cb-722", null));
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
"a presence contact that lands before worktree registration must still reach READY");
}
@Test
void worktreeArchitectAcquireAlsoBindsItsSlot() {
FakeHerdr herdr = new FakeHerdr();