Compare commits

..

42 Commits

Author SHA1 Message Date
Dai Ha c043d149cf fleetd #775: name the trigger flag for what it now reads
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 1m0s
CI / build (push) Failing after 1m57s
The flag tested leader.tab() and was named for it. It now tests whether
fleet.leaders has any entry at all, so the old name states a condition the
code no longer checks.
2026-10-05 14:21:19 +02:00
Dai Ha fc786d0f67 fleetd #775: say which remedy applies to a lead vs a collaborator
The ticket's correction comment pointed out the refusal message still
implied removing a lead's tab helps, when only placement: tab does.
Restate the message so the lead and collaborator remedies are not
conflated, and keep the variable name the correction specified.
2026-10-05 14:19:17 +02:00
Dai Ha b314cb4d51 fleetd #775: pane-placement guard trigger keys on lead existence, not tab:
The guard's lead half now fires whenever fleet.leaders has any entry,
since a lead's tab is always labelled by the fixed LEAD_TAB_LABEL
constant regardless of its own deprecated tab: field. The refusal
message no longer advises removing a lead's tab:, which cannot
satisfy the guard any more.
2026-10-05 14:19:17 +02:00
Dai Ha a3d296f639 fleetd #770: the lead identity check is the label AND the space
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 2m11s
The redeploy skill's check 4 and the script's closing hint both told the
operator that a lead is found by fleet.leaders.*.tab. Identity is now the
fixed 'lead' label together with the lead's configured workspace, so both
would have sent a reader to a key that no longer decides anything.
2026-10-05 14:01:54 +02:00
Dai Ha 79423787d0 fleetd #770: lead identity keys on the space, not the tab label
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 1m57s
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 53s
CI / build (push) Failing after 1m52s
The lead tab label becomes a fixed constant (Leader.LEAD_TAB_LABEL = "lead");
fleet.leaders.<name>.tab is now optional legacy, matched case-insensitively
alongside the constant via Leader.acceptedLabels(). The uniqueness boundary
between leads moves from the exact tab text to the workspace: FleetConfig
refuses two leaders that share a workspace, LeadTabScanner indexes lead
labels per space (collaborators stay space-agnostic), and
LeadLauncher.leadNameOf/countLeads require both the accepted label and the
lead's own space to match, so a legacy-labelled tab in the wrong space never
counts and a daemon restart never double-spawns a second lead next to a live
one. Config validation also refuses a fleet.tabLabel template or a
collaborator tab that can render as the fixed lead label.
2026-10-05 13:49:14 +02:00
Dai Ha 364b229db9 fleetd #771 step 1: report each pane's workspace label in fleet_list
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 1m57s
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m50s
Adds workspaceLabel next to workspaceId on fleet_list's panes row, read
from herdr's workspace.list via a new PaneLocator.workspaceLabelsByWorkspaceId().
An observer's reduced row still excludes it; a herdr failure or an unknown
workspaceId yields workspaceLabel: null without costing the rest of fleet_list.
2026-10-05 11:37:41 +02:00
Dai Ha ac436efefb fleetd #761: invariant 4 — a lead's own pane needs an empty input box
CI / shell-tests (push) Failing after 11s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m58s
2026-10-05 11:17:33 +02:00
Dai Ha 7dad054045 Merge remote-tracking branch 'origin/worker/761-draft-aware-lead-nudge-d05850-4' into vfy/761
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 1m0s
CI / build (push) Failing after 2m16s
2026-10-05 11:09:40 +02:00
Dai Ha 4e414c475f fleetd #761: read the input box marker the live TUI actually draws
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 2m1s
The gate matched the box line as "│ >", which appears zero times on a current
Claude Code pane. EMPTY was unreachable, so every lead nudge was held forever.

Match the box as a line *starting* with "❯" or with "│ >", keeping the older
bordered layout readable. A marker further along a line is transcript text — a
caret the operator quoted — so it no longer counts, and the last matching line
is still the live box because the detection region carries scrollback above it.

Look for the generating marker only from the box line down, for the same
reason: an earlier turn's "esc to interrupt" survives in that scrollback, and
holding on it would be the same unreachable-EMPTY failure by another route.

Fixtures: add IDLE_PROMPT_CARET / DRAFTED_PROMPT_CARET from a live pane and
point every lead-pane fake at them. The bordered constants stay, now covering
the older client. Against the old marker these fixtures fail 55 tests across
PromptBoxTest and the three loop tests, which is the production bug reproduced.

cd fleetd && mvn clean install -> BUILD SUCCESS, MVN_EXIT=0,
Tests run: 2164, Failures: 0, Errors: 0, Skipped: 0 (179 surefire XML files).
2026-10-05 11:07:40 +02:00
Dai Ha 9084667493 fleetd #761: hold a lead nudge while the operator's prompt box holds a draft
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Failing after 2m13s
herdr's agent.prompt pastes AND submits in one call, so a nudge arriving
while the operator is mid-sentence submitted their unfinished line with the
nudge glued to it. AgentStatus.injectable() cannot see this: it describes
the agent, and an idle agent reports the same status whether its input box
is empty or holds a half-typed line.

New herdr/PromptBox reads the pane's `detection` region — the same region
StatusRefiner uses, and the one the input box is drawn in — and clears a
delivery only when the box is positively empty. A box with characters, a
pane it cannot recognise, and a failed read all hold, because a held nudge
is recoverable and a submitted half-line is not. Whitespace and a cursor
block count as empty. After 20 consecutive holds for one target it logs one
warning, so a box that never clears is visible rather than silent; the
warning repeats only after the box has cleared again.

Wired into the three paths that nudge a lead's own pane:
  - ReplyPushLoop.decide -> WAIT_BUSY (the pending work is re-read next tick)
  - LeadHeartbeatLoop.tick -> a new DRAFT_HELD action that spends neither the
    quiet budget nor the one context notice per HIGH stretch
  - LeadCoordLoop.tick -> the peer message stays held and unacked

Not wired into inject/Injector: no human types into a spawned member's pane,
so it would buy nothing and cost a herdr agent.read per member poll. The
heartbeat reads the pane only for a tick that would otherwise send.

Tests. FakeHerdr gains detectionText(), because one readText cannot be both
a worker's transcript (what the completion scrape reads) and a lead's empty
prompt; it falls back to readText so no existing fixture changes meaning.
Fixtures for the lead-nudge paths now state what their pane shows, since the
behaviour depends on it. 11 new behavioural tests across the three loops plus
InjectorTest, and 13 for the classifier; all 10 loop tests were run against
the unpatched loops first and fail there.

mvn clean install: BUILD SUCCESS, Tests run: 2161, Failures: 0, Errors: 0,
Skipped: 0 (summed from target/surefire-reports).
2026-10-05 10:49:50 +02:00
Dai Ha 2289e94223 fleetd #759: fix the fleet_reply comment and the hand-copied role list
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Failing after 1m59s
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 1m59s
Finding 4: the comment above fleet_reply's handler claimed the authz check
asks whether the caller is a worker at all. It actually checks terminal
ownership (Authz.java REPLY/ASK -> caller.ownsSession), which is why an
observer can reply on its own pane with no role test involved.

Finding 5: fleet_list's tool description hardcoded 'architect/dev/reviewer',
missing hunter. Added MemberRole.wireNames() (pulled out of parse()'s error
message builder, which now calls it too) and used it in the description so
the list can't drift again.
2026-10-05 10:44:55 +02:00
Dai Ha 92adfcfae5 fleetd #756/#758: the canonical block no longer says panes are unlistable
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 2m13s
The table row for an unconfigured pane told every session "Neither
fleet_list nor ListAgents lists these". PR #762 made that false for
fleet_list, and a stale note of this shape is the worst kind: it tells a
future session it cannot do the thing at the moment doing it is the job.

Two edits:

- The observer definition now says where an observer finds a target id,
  which is the one thing it could not learn before.
- The table row names the panes array, says the row is full for a lead, an
  architect or a collaborator and filtered for an observer, and keeps the
  herdr tab list join as the fallback. ListAgents still lists none of them.

wiki/7-Use-Cases.md regenerated from this block in the wiki submodule at
4872227; the sync check prints in sync: True.
2026-10-05 10:39:32 +02:00
Dai Ha 5f7f388e69 Merge remote-tracking branch 'origin/worker/756-758-observer-pane-discovery-7e6ffd-1' into vfy/762
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m56s
2026-10-05 10:29:49 +02:00
Dai Ha 39accf73e6 fleetd #759: fix the third copy of the role claim, and drop two references that rot
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m58s
ConnectionIdentity's Caller record carried the same wrong rule as the two
places PR #760 fixed: it called the terminal a worker's, and read a null
terminal as the primary. The brief for #760 named only two of the three spots.

MemberPresence pointed at FleetMcp.markTrackedCallerPresent by name inside
{@code}, which the compiler does not check, and inject has no dependency on
mcp so a {@link} would add a cross-package reference. States the principle
instead, which stays true whichever roles qualify. Drops a ticket key.
2026-10-05 10:26:36 +02:00
Dai Ha 5c2f296bc3 fleetd #756/#758: panes array reads the architect-slot gate and widens to observer
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 2m9s
paneRole now reads CallerResolver#boundToArchitectSlot (made public, no
second definition) so a slot-bound pane with no live member reports
"architect", matching what sendableObserverTarget already allowed as a
SEND target.

panesVisibleTo now admits an observer, since an observer holds SEND to
another observer pane. Its panes rows are filtered to
sendableObserverTarget and reduced to sessionId/label/status/role/
deliverable; every other caller's rows are unchanged.
2026-10-05 10:24:00 +02:00
Dai Ha 0f2ec7a6b5 fleetd #759: fix three role-model comments that describe the pre-CB-501 rule
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 2m7s
Role.PRIMARY, ConnectionIdentity (class javadoc + callerTerminal), and
MemberPresence's class javadoc each state a role model the code no longer
implements. Comment-only change.
2026-10-05 10:18:01 +02:00
Dai Ha 02eff4c532 fleetd #743: an observer may now send to another observer
CI / shell-tests (push) Failing after 13s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 2m4s
Three claims in the canonical block went false when the observer SEND grant
merged, and this file is the instruction surface the bridge ships:

  - the observer role definition said "never SEND"
  - invariant 3 listed send as "lead, architect, or collaborator"
  - the hand-opened-pane row said such a pane "cannot fleet_send back"

The grant is narrow. Authz permits an observer's SEND only when
CallerResolver.sendableObserverTarget() classifies the target as an observer
too, so a lead, a collaborator, an architect slot and a spawned member are
each refused. TASK_READ stays denied, so #705 remains closed.

Measured on the merged tree: mvn clean install exit 0, 2137 tests from 178
surefire files, 0 failures. Mutating away the grant's call site
(FleetMcp.java:474) kills exactly FleetMcpObserverSendDeliveryTest, so the
behaviour is pinned and not only the predicate.
2026-10-05 09:42:43 +02:00
Dai Ha 92b6d9406b Merge remote-tracking branch 'origin/worker/743-observer-send-4706db-6' 2026-10-05 09:37:14 +02:00
Dai Ha c4498607e1 Merge remote-tracking branch 'origin/worker/743-pane-discovery-ad5b75-5' 2026-10-05 09:37:14 +02:00
Dai Ha e42eab5b4c fleetd #743: pin GET /agents' tab-label merge with a positive test
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 2m5s
2026-10-05 09:34:15 +02:00
Dai Ha b1d2cb48ac fleetd #743: make the pane-label scan best-effort, trim justification comments
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m56s
A workspace.list/tab.list failure in the label scan no longer costs the
caller the agent roster (GET /agents) or the leads/members/capacity/
coordinator rows (fleet_list) that never needed it. Both call sites now
fall back to an empty label map on HerdrException, so a pane row still
renders with label:null instead of the whole response failing.

Also cuts four comments down to the current contract, per the project's
comment rule: dropped the reviewer-facing justification from Fleetd's
deliverableTo javadoc, panesVisibleTo's javadoc, the fleet_list handler's
inline comment, and the panesVisible assembly-gate comment, and removed
the two fragments describing what a test must do.
2026-10-05 09:19:16 +02:00
Dai Ha 001367d82c fleetd #743: drop the redundant source-scrape test for attributeIfObserver
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Failing after 2m4s
FleetMcpObserverSendDeliveryTest already kills the same mutation end to
end (it asserts the exact attributed text herdr receives), and the code
quality rule in CLAUDE.md caps new source-text tests in this file at the
existing count.
2026-10-05 09:16:29 +02:00
Dai Ha 9fdcaaa8fd fleetd #743: let one observer pane message another, and nothing else
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 56s
CI / build (pull_request) Failing after 2m0s
Authz.SEND now grants an observer a narrow path: it may reach only a
target that CallerResolver.sendableObserverTarget() would itself
resolve as OBSERVER, never a lead, a collaborator, or a live spawned
member's terminal. This mirrors the existing collaborator SEND clause
rather than adding an unconditional caller.isObserver() grant, which
fleetd #705 already rejected as too broad.

Because the receiving pane cannot otherwise tell an observer's SEND
apart from a human paste, FleetMcp.attributeIfObserver prefixes the
delivered text with the sender's own daemon-resolved terminal on both
the MCP and REST entry paths, for exactly this one new path.

FleetAppAuthTest's start() helper never wired a spawnedMemberRole, so
its "worker" fixture actually resolved as OBSERVER under the real
CallerResolver -- invisible before because OBSERVER and WORKER shared
the same (zero) SEND grant. Granting OBSERVER a real SEND surfaced it:
two SEND-denial tests started passing for the wrong caller. Fixed the
fixture to resolve term_a as a live DEV, matching the helper's own
documented contract.
2026-10-05 09:03:16 +02:00
Dai Ha 459a523e2c fleetd #743: expose herdr pane discovery (tab labels) over REST and MCP
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m0s
CI / build (pull_request) Failing after 2m15s
GET /agents now carries each agent's tab label, merged in from the same
herdr daemon(s) the roster is drawn from. fleet_list gains a panes array
with the same label, the sessionId fleet_send takes as a target, the
role the daemon resolves that pane as, and the deliverable gate the
injector itself enforces (Fleetd#deliverableTo, now public). Gated like
leads/collaborators (primary, architect, collaborator), not bare READ,
since a tab label and a member's cwd are not roster facts every READ
caller may see.
2026-10-05 08:58:58 +02:00
Dai Ha 291dc02c77 fleetd #743: document that a hand-opened pane is reachable with no config
The lead's intent->tool table had no row for messaging an unconfigured pane,
and the user-scope instruction file said outright that the fleet has no route
to an interactive session unless an operator registers it as a collaborator.
That claim is false and it is load-bearing: a session reading it concludes the
exchange is impossible and stops, which is what happened here.

Delivery is gated on presence, not on SEND. contextExtractor runs on every MCP
request including initialize, markTrackedCallerPresent enrols an observer into
MemberPresence, and deliverableTo tests presence before the lead and
collaborator maps. So connecting the server is the enrolment, and fleet_reply
is gated on owning your own pane, which every pane does.

Add the table row, and add the enrolment side of the deliverability gate to the
flows page next to the existing "a spawned member is not deliverable until it
has mounted the MCP" bullet, which is the same gate read the other way.

Measured on a real pane, not a fake: trinotes answered with no fleet config, no
restart, and its fleet_* tools still deferred and unloaded.
2026-10-05 08:56:01 +02:00
Dai Ha be835aa259 Merge branch 'worker/749-edge-baseline-28d1a0-3'
CI / shell-tests (push) Failing after 12s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 2m4s
2026-10-05 07:47:18 +02:00
Dai Ha 80506c79d0 fleetd #748: drop a cross-reference the comment cleanup left dangling
errorPatternCoverageLine pointed at exhaustedPatternCoverageLine's javadoc
"for the measured swap mutation this pairing guards against". That narrative
was removed from the destination, so the pointer led nowhere.

The pairing with BUILT_IN_DEFAULT is the contract and it stays. What the
mutation proved belongs in history, not in the comment.
2026-10-05 07:46:59 +02:00
Dai Ha 6677ec8c63 fleetd #749: pin PackageCyclesTest's exceptions to exact edges, not whole packages
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m8s
CI / build (pull_request) Failing after 2m25s
ignoreCycle() used to exempt every dependency between two packages, in both
directions, for the whole package. That meant a brand new dependency added
later between an already-excepted pair (auth/mcp, mcp/msg, inject/msg,
metrics/msg, msg/session) was silently exempted too, exactly where the
msg package makes the gate matter most.

Replace the package-wide ignore with a frozen baseline of the 45 exact
origin-class -> target-class edges that exist today between those five
pairs, and ignore only those via SliceRule.ignoreDependency(String, String).
A new dependency between a baselined pair is not in that set, so it is no
longer ignored and the existing beFreeOfCycles() check (or, when the new
edge alone would not form a cycle, a dedicated set-equality check) fails
and names the exact origin class, target class and package pair.

The set-equality check also fails on a baseline entry whose dependency no
longer exists in the code, so a removed edge cannot rot in the baseline
and mask the pair's eligibility for the ticket #131 removal steps. Rewrote
the javadoc to describe only the current contract.
2026-10-05 07:45:51 +02:00
Dai Ha ca0c965932 Merge branch 'worker/748-dead-comment-refs-f42ac5-4' 2026-10-05 07:44:44 +02:00
Dai Ha e8ab933cbc fleetd #748: fix dead test-class references and orphaned javadoc blocks
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 2m1s
Drop 2 dead *WiringTest names from 3 comment sites (renamed to
*AssemblyTest), rewriting each sentence to state the code's guarantee
instead of naming a test class. Reattach 5 javadoc blocks that were
orphaned behind a second /** block to the member they actually
describe, trimming history/evidence text down to the current contract
per the project's comment rule. Comment-only; no production logic
changed.
2026-10-05 07:42:55 +02:00
Dai Ha 2adb950a12 fleetd #748: five code-quality rules enter the repo
Two architects reviewed the codebase independently on separate backends and
both returned MIXED, not "a mess": 0 public mutable fields, 12 extends of
which 8 are exceptions, 120 records, and 2 files over 1000 code lines rather
than the 9 a total-line count suggests. Both rejected a Clean Code section,
a SOLID list, a pattern catalogue, class and method length limits, and a
coverage gate as text that would change no behaviour.

The comment rule they were briefed against was never in this repo. It lives
in the operator's personal global config, so it was not in git, never
reviewed, and not versioned with the code. Its "goes in the ADR" line was
unfollowable because this project has no ADR, which sent load-bearing
knowledge to a destination that does not exist. Rule 2 redirects it to
docs/<subject>.md, which exists.

Rule 4 covers what neither architect ranked first and both described: a
testing problem relieved by reshaping production code. 54 static factories
on Fleetd, 7 volatile race hooks in MessageService, 8 tests asserting on
main source as text, and a 452-line composition root across 8 tickets.

Rule 4's judgement half is marked as having no mechanism. A cap stops a
count growing; no test distinguishes a good decomposition from a bad one.

Verified: canonical block still byte-identical with wiki/7-Use-Cases.md.
2026-10-05 07:36:20 +02:00
Dai Ha 7f0c4a8464 fleetd #737: the handover skill carries the outstanding tickets forward
CI / shell-tests (push) Failing after 12s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 2m7s
fleet_handover{open} now returns outstandingTickets and openAsks. The
successor keeps the authority to poll and answer them but not the ids, and
an uncollected terminal ticket loses its reply at the ticket TTL.
2026-10-05 06:09:26 +02:00
Dai Ha 682991a846 Merge branch 'worker/737-9c61d3-4'
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 50s
CI / build (push) Failing after 1m52s
2026-10-05 06:07:33 +02:00
Dai Ha abe617c48c Merge branch 'worker/737-a263f3-3'
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m3s
CI / build (push) Failing after 2m20s
2026-10-05 06:01:49 +02:00
Dai Ha 886ce1521d Merge remote-tracking branch 'origin/main' into worker/737-9c61d3-4
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m59s
2026-10-05 06:00:49 +02:00
Dai Ha d41aff4012 fleetd #737 unit 5 correction: outstanding() reports terminal-phase tickets too
A DONE/FAILED ticket nobody has polled yet is destroyed on a timer by
pruneTerminalTickets' completion-based TTL, while a PENDING ticket is not
going anywhere. Drop the !future.isDone() filter so outstanding() reports
every ticket the caller owns that is still in tasks, with its real phase.
2026-10-05 06:00:43 +02:00
Dai Ha 8a1d73b39e fleetd #737 unit 4: key the rollover single-flight claim on the lead's name
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 1m59s
rollingByTerminal keyed the single-flight lock on p.leadTerminal(), the pane
address. A roll replaces the pane, so a second roll of the same lead opened
from the new terminal landed on a different map key and could run concurrent
with the first roll's still-in-flight continuation.

PendingRollover now carries rolloverKey, resolved once in open() from
leadNameForTerminal while the lead is certainly still live, falling back to
the terminal itself when the name resolves null or blank. confirm()'s claim
and both release sites (the continuationRunner-rejection catch and
runRollover's finally) use the carried key instead of recomputing it, since
leadNameForTerminal no longer resolves the old terminal by release time.
NOT_YOUR_ROLLOVER and the pane-teardown calls stay keyed on p.leadTerminal(),
unchanged.

rollingByTerminal is renamed rollingByLead to match.
2026-10-05 05:57:59 +02:00
Dai Ha 70a735b638 fleetd #737 unit 5: fleet_handover{open} reports outstanding tickets and open asks
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 2m10s
MessageService.outstanding(callerOwner) lists the caller's non-terminal
delegations and the subset paused in fleet_ask, filtered by the same
ownsTicket rule poll() already uses. FleetMcp threads the caller's owner
key into handover/handoverOpen and adds outstandingTickets/openAsks to
the open() JSON, alongside the unchanged token/handoverPath/requestedAtMillis.
2026-10-05 05:56:06 +02:00
Dai Ha 803c91ea6c fleetd #737 unit 3: name the resolving accessor in LeadRollover's javadoc
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m0s
CI / build (push) Failing after 2m7s
The heartbeat loop now resolves its nudge target through
PrimaryRegistry#currentPrimaryTerminal(), so the sentence describing what a
background loop with no caller uses named the raw accessor instead.
2026-10-05 05:37:32 +02:00
Dai Ha d438a74575 Merge branch 'worker/737-20d1d9-1' 2026-10-05 05:37:32 +02:00
Dai Ha 7467ffa252 fleetd #737 unit 6: drop the stale unnamed-primary claim from the REST status comment
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 2m9s
The comment described the pre-#737 rule, where a null caller key read every
ticket. The pendingAsk gate now matches the unnamed primary's key like any
other, so the exemption it named no longer exists.
2026-10-05 05:32:59 +02:00
Dai Ha 6ab3a81af7 fleetd #737 unit 6: stop the unnamed primary sharing the internal bypass
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m59s
ownsTicket treated a null callerOwner as "read everything", conflating
the internal no-check test seam with a real unnamed primary's owner
key. Give them two different values: a named INTERNAL_NO_OWNER_CHECK
marker for the test-only bypass, and null as just another owner key
that must equal the ticket's recorded creatorOwner (including a
null-to-null match, so an unnamed primary still owns its own tickets).
Applies to both ownsTicket call sites, poll and pendingAsk.
2026-10-05 05:24:37 +02:00
55 changed files with 3186 additions and 563 deletions
+11 -2
View File
@@ -133,8 +133,17 @@ A handover file is a record of state and decisions. It is not a diary.
**Run the three steps in this order. The order is not a style choice — the wrong order is
refused.**
1. **`fleet_handover{action: "open", reason: "<why now>"}`.** It returns a `token` and the
`handoverPath` you must write to. Nothing has happened to your pane yet.
1. **`fleet_handover{action: "open", reason: "<why now>"}`.** It returns a `token`, the
`handoverPath` you must write to, and `outstandingTickets` plus `openAsks`. Nothing has happened
to your pane yet.
**Copy `outstandingTickets` and `openAsks` into the handover file.** Your successor keeps the
authority to poll those tickets and answer those asks, because both are gated on the lead's
name, which does not change when your pane does. What it does not keep is the ids — they exist
only in your context and in this response. A ticket already in a terminal phase is the urgent
one: its reply lives only in memory and is deleted once the ticket TTL passes, so an uncollected
report is lost for good. Poll those before you confirm, or name them in the file so your
successor polls them first.
**Write to exactly that path, and do not resolve it yourself.** It is always absolute, even when
the operator configured a relative `handoverPath`: fleetd resolves a relative one against your
+6 -3
View File
@@ -59,9 +59,12 @@ if the script is unavailable or a step fails, this is what it was protecting you
3. **A restart is the only way deferred config keys take effect.** That is usually the reason to do
it. The startup log names which keys it accepted and which it deferred — read those lines rather
than assuming.
4. **Re-check identity afterwards.** Call `fleet_whoami` and confirm it still answers `primary`. The
lead is found by its tab label (`fleet.leaders.*.tab`), and a lead whose tab no longer matches is
demoted to worker, which refuses every orchestration call.
4. **Re-check identity afterwards.** Call `fleet_whoami` and confirm it still answers `primary`. A
lead is found by two things together: its tab is labelled `lead`, and that tab sits in the space
named by `fleet.leaders.<name>.workspace`. Both must match, so a renamed tab *and* a space whose
label differs from the config each demote the lead to worker, which refuses every orchestration
call. A `tab:` still in config is accepted as a second label for that lead, and the daemon logs
one deprecation warning naming it at startup.
5. **Prove the new jar is the one running.** Confirm a *fresh* `fleetd listening` line at the end of
`fleetd/fleetd.out`, dated after the restart. An old daemon that never died looks identical from
the outside.
+47 -5
View File
@@ -33,8 +33,11 @@ gate uses. A worker also carries its `sessionId`, `profile`, `worktree` and `bra
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.
to `READ`/`METRICS`, to `REPLY`/`ASK` on its own pane, and to `SEND` only to a target that resolves
as an observer too — never to a lead, a collaborator, or a spawned member, and never a ticket. It
finds such a target in `fleet_list`'s `panes` array, which for an observer is filtered to exactly
what it may send to and reduced to `sessionId`, `label`, `status`, `role` and `deliverable`.
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
@@ -62,14 +65,22 @@ and the sender silently receives nothing. Fail toward the recoverable error.
2. **The bridge is the only channel.** Text you print in your terminal reaches nobody — the other
side cannot see your screen. An answer that isn't in a `fleet_*` call is silently discarded.
3. **Identity comes from the connection, never an argument.** Workers never pass a target; you
cannot act as another session. Spawn/stop/drain are lead-only; **send is lead, architect, or
collaborator** — and a collaborator may send only to a lead or another collaborator, never to a
spawned member's terminal; reply/ask are only-as-itself — any peer may answer for its own pane,
cannot act as another session. Spawn/stop/drain are lead-only; **send is lead, architect,
collaborator, or observer** — and a collaborator may send only to a lead or another collaborator,
never to a spawned member's terminal, while an observer may send only to another observer pane;
reply/ask are only-as-itself — any peer may answer for its own pane,
and for no other. A call outside your role is refused, not queued.
4. **Delivery is status-gated: one message per turn.** Don't busy-poll a peer's terminal and don't
re-send because a call looks slow — the bridge delivers when the peer is `idle`, `blocked` or
`done`. A spawned member must **also** have mounted the bridge MCP: until it has, it is not
deliverable, and a send waits on that gate for ~60s and then fails without ever reaching its pane.
**A lead's own pane has a second gate: its input box must be empty.** The multiplexer pastes and
submits in one step, so a delivery that lands while the operator is typing submits their
half-written line. A heartbeat, a ticket nudge and lead-to-lead mail therefore wait until the box
is clear, and a pane the daemon cannot read as a box waits too. Nothing is lost — every one of
those paths retries — but a lead that leaves text sitting in its box receives nothing until it
clears, and the only sign is one warning in `fleetd.out` after 20 held checks in a row. Delivery
to a *member* is not gated this way, because nobody types in a member's pane.
5. **Never move a fleet session, pane or peer except through the bridge.** The bridge owns policy;
the multiplexer owns PTYs. Any route that changes fleet state without the bridge's checks
bypasses every rule above — the `herdr` CLI and its socket are the usual example.
@@ -178,6 +189,7 @@ you decide.
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` reports a `collaborators` array, and each row carries that peer's `name` and the `sessionId` you send to. It is visible to you, to an architect and to another collaborator, never to a worker. Coordination only, **never** a task |
| Message an **unconfigured pane** — a tab a person opened by hand | `fleet_send{sessionId: <their terminal>, content}` — it needs **no** `fleet.collaborators` entry and no restart, because a pane becomes deliverable the moment its agent connects the bridge MCP. `fleet_list`'s `panes` array reports every such pane with its label and the terminal id to send to — the full row for you, an architect or a collaborator; filtered and reduced for an observer. **`ListAgents` still never lists these**, and joining `herdr tab list` to `GET /agents` on `tab_id` stays the read-only fallback if the array is missing. Such a pane resolves as an `observer`: it can answer you with `fleet_reply`, and it can `fleet_send` to another observer pane, but never to you. Coordination only, **never** a task |
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
| Read your own held lead-to-lead mail (no ack) | `fleet_poll{coordId: <your own coord-id, from fleet_list's coordinator.selfId>}` — primary-only; never acks, so `fleet_list`'s `held[]` still shows it after. `fleet_list`'s `held[]` gives only a truncated preview — this is the only way to read the full body |
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
@@ -502,3 +514,33 @@ to replace them.
Prefer the unnamed lambda parameter `_` for required-but-unused params; a non-public
`static void main(String[])` is valid (JEP 512) and boots via `java -jar`.
## Code quality — five rules, and what each already cost (enforced)
Measured at `7f0c4a8`: 124 main files, 36,278 lines, of which **17,850 are code** — 44% is comment,
and only **two** files exceed 1000 *code* lines. Encapsulation and inheritance are already sound (0
public mutable fields; 12 `extends`, 8 of them exceptions; 120 records). So there is **no Clean Code
section, no SOLID list and no pattern catalogue** here: two architect reviews rejected those
independently as text that would change no behaviour. These five rules are the whole standard.
1. **A comment states the current contract or a current maintainer constraint — nothing else.** No
tickets, history, dates, measurements or review rationale; those go in the commit message or the
MR description. Source code only — this rule never applies to Markdown.
2. **A javadoc block stops at 30 lines.** Longer means it is a design argument, so it moves to
`docs/<subject>.md` and is linked in one line. The longest here is 235 lines
(`config/ConfigRef.java`) and the knowledge in it is load-bearing: **move it, never delete it.**
This project has **no ADR** — subject pages under `docs/` are the destination.
3. **A comment in main source never names a test class.** There are 44 such names in 76 places and
**2 are already dead**, because a name inside `{@code}` is invisible to the compiler and rots in
silence. Say what the code guarantees; the test is found by looking.
4. **Never relieve a testing problem by reshaping production code.** `Fleetd` carries 54 static
factories, `MessageService` carries 7 `volatile` race hooks, and 8 tests assert on main source as
*text*. Make the part injectable instead. `FleetdAssembly.assembleAndStart` is 452 lines and may
not grow; no new source-text test may be added.
5. **No new package cycle, and no widening of a recorded one.** Five pairs are frozen as an exact
edge baseline in `PackageCyclesTest` — four of them involve `msg`.
Rules 1, 2, 3 and 5 have build checks, and Gitea CI runs them on every PR, so they bind members too.
**Rule 4's judgement half has no mechanism**: a cap stops a count growing, but no test tells a good
decomposition from a bad one. That half is a review obligation, and saying so is deliberate — a rule
dressed as a gate it does not have is worse than an honest review item.
+6
View File
@@ -182,6 +182,12 @@ Two consequences a lead feels directly:
is fine; the message simply waits, and then restarts the member when it next goes idle.
- **A spawned member is not deliverable until it has mounted the MCP.** Until then a send waits on
that gate for about 60 seconds and then fails without ever reaching the pane.
- **The same gate is what makes an unconfigured pane deliverable.** `contextExtractor` runs on
every MCP request, `initialize` included, and `markTrackedCallerPresent` enrols a spawned member
*or* an observer into `MemberPresence`; `deliverableTo` then tests presence before the lead and
collaborator maps. So mounting the server is the enrolment, and a tab a person opened by hand can
be sent to with no config and no restart. It answers with `fleet_reply` — it cannot `fleet_send`,
because `Authz` keeps `SEND` to a primary, an architect or a collaborator.
`UNKNOWN` is deliberately neither injectable nor a pickup. A pane whose status cannot be read is
not a pane that is safe to write to — see fleetd #176 for what happens when a gate treats an
+13 -26
View File
@@ -234,7 +234,7 @@ public final class Fleetd {
* <p>Both sets are read through their supplier on each call rather than snapshotted, so a lead or
* collaborator discovered by {@code leadScan} after startup becomes deliverable without a restart.
*/
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads,
public static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads,
Supplier<Map<String, String>> collaborators) {
return target -> presence.isPresent(target) || leads.get().containsKey(target)
|| collaborators.get().containsKey(target);
@@ -336,22 +336,6 @@ public final class Fleetd {
}, reasonByCredential::get);
}
/**
* fleetd #415 (review follow-up): package-private factory for the CB-578 stage A {@code
* exhaustedPattern} startup coverage line, paired explicitly with {@link
* CompletionResolver.UnsetMeaning#OFF} — {@code exhaustedPattern} has no fallback, so a
* profile with none configured really does have the classification off.
*
* <p>Extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
* #worktreeBranchLookup} were: {@code coverage()}'s own tests ({@code CompletionResolverTest})
* prove it words {@code OFF} and {@link CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT}
* correctly when a test supplies the meaning itself — they cannot prove {@code main} pairs the
* right meaning with the right key, which is the actual fleetd #415 defect. <b>Measured:</b>
* swapping the {@code UnsetMeaning} arguments between this method and {@link
* #errorPatternCoverageLine} — recreating #415's defect with the two keys exchanged — compiled
* with 0 errors and left all 1506 existing tests green before {@code
* FleetdPatternCoverageLineTest} was added to catch exactly that swap.
*/
/**
* fleetd #446 follow-up: the criterion-2 WARNING text — "name the fix, not just the fact" —
* for a profile whose {@code model:} is configured. Extracted out of the {@code
@@ -386,19 +370,22 @@ public final class Fleetd {
+ "s quarantine above is the only thing keeping new spawns off it for now";
}
/**
* Package-private factory for the {@code exhaustedPattern} startup coverage line, paired
* explicitly with {@link CompletionResolver.UnsetMeaning#OFF} — {@code exhaustedPattern} has
* no fallback, so a profile with none configured really does have the classification off.
*/
static String exhaustedPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
return CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
allProfiles, configuredProfiles);
}
/**
* fleetd #415 (review follow-up): the {@code errorPattern} counterpart of {@link
* #exhaustedPatternCoverageLine}, paired explicitly with {@link
* CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT} — an unset {@code errorPattern} still runs
* backend-error classification against {@code CompletionResolver}'s built-in {@code
* BACKEND_ERROR} pattern, so the empty case is not "off". See {@link
* #exhaustedPatternCoverageLine}'s javadoc for the measured swap mutation this pairing guards
* against.
* The {@code errorPattern} counterpart of {@link #exhaustedPatternCoverageLine}, paired
* explicitly with {@link CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT} — an unset
* {@code errorPattern} still runs backend-error classification against
* {@code CompletionResolver}'s built-in {@code BACKEND_ERROR} pattern, so the empty case is
* not "off".
*/
static String errorPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
return CompletionResolver.coverage("errorPattern", CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT,
@@ -734,8 +721,8 @@ public final class Fleetd {
* {@code CompletionResolver} constructor call — provably untested wiring, the whole reason
* fleetd #248 exists: dropping that one argument (passing {@code _ -> null} instead) compiled
* clean and left every test green. Extracted here, {@code main} now calls this factory instead
* of building the lambda inline, and a source assertion on that call site
* ({@code FleetdCompletionResolverWiringTest}) proves the argument is still actually passed.
* of building the lambda inline, so the argument reaching the {@code CompletionResolver}
* constructor is a named, directly testable call rather than an inline lambda.
*
* <p>Takes the roster as a plain {@link Supplier} — not a {@link SessionManager} — so this is
* directly testable with a hand-built session list; no real {@code SessionManager} (launcher,
@@ -254,18 +254,28 @@ final class FleetdAssembly {
if (leadTerminals.size() > 1) {
log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values());
}
// CB-531/CB-579: discover leads by the tab labels the operator writes, one scanner per
// configured lead's own exact `tab:` label. fleetd #669: the same scan also recognises a
// configured collaborator's tab, so one herdr pass answers both.
// Discover leads by their tab labels, one scanner per configured lead's own space. fleetd
// #669: the same scan also recognises a configured collaborator's tab, so one herdr pass
// answers both.
final Supplier<Map<String, String>> leads;
final Supplier<Map<String, String>> collaboratorTerminals;
var leaders = cfg.fleet().leaders();
var collaboratorsConfig = cfg.fleet().collaborators();
if (!leaders.isEmpty() || !collaboratorsConfig.isEmpty()) {
Map<String, String> tabToName = new LinkedHashMap<>();
Map<String, Map<String, String>> leadLabelsBySpace = new LinkedHashMap<>();
Map<String, String> spaceByLeadName = new LinkedHashMap<>();
leaders.forEach((name, leader) -> {
if (leader != null && leader.tab() != null && !leader.tab().isBlank()) {
tabToName.put(leader.tab(), name);
if (leader == null) {
return;
}
spaceByLeadName.put(name, leader.workspace());
Map<String, String> labelsHere = leadLabelsBySpace
.computeIfAbsent(leader.workspace(), k -> new LinkedHashMap<>());
leader.acceptedLabels().forEach(label -> labelsHere.put(label, name));
if (leader.tab() != null && !leader.tab().isBlank()) {
log.warn("lead '{}' (fleet.leaders.{}) still configures tab: \"{}\" — deprecated, "
+ "the lead tab label is now fixed to '{}'",
name, name, leader.tab(), FleetConfig.Leader.LEAD_TAB_LABEL);
}
});
Map<String, String> collaboratorTabToName = new LinkedHashMap<>();
@@ -283,13 +293,13 @@ final class FleetdAssembly {
? 10
: leaders.values().iterator().next().scanIntervalSeconds();
// This must use the lead daemon: scanning member tabs would demote the lead to a worker.
LeadTabScanner scanner = new LeadTabScanner(herdr, tabToName, collaboratorTabToName, Set.of(),
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), ports.nanoClock());
LeadTabScanner scanner = new LeadTabScanner(herdr, leadLabelsBySpace, collaboratorTabToName,
Set.of(), TimeUnit.SECONDS.toNanos(scanIntervalSeconds), ports.nanoClock());
leads = scanner;
collaboratorTerminals = scanner::collaborators;
log.info("lead/collaborator scan: tabs {} host a lead, tabs {} host a collaborator "
+ "(rescan every {}s, shared fleet space)",
tabToName.keySet(), collaboratorTabToName.keySet(), scanIntervalSeconds);
log.info("lead/collaborator scan: space per lead {}, tabs {} host a collaborator "
+ "(rescan every {}s)",
spaceByLeadName, collaboratorTabToName.keySet(), scanIntervalSeconds);
} else {
leads = () -> leadTerminals;
collaboratorTerminals = Map::of;
@@ -70,33 +70,58 @@ public final class Authz {
*/
public static final Predicate<String> NO_KNOWN_LEAD_OR_COLLABORATOR = target -> false;
/**
* The fail-closed classifier for an observer's {@code SEND}: answers no for every target, so
* the grant is refused unless a caller supplies a real one. {@code
* CallerResolver#sendableObserverTarget()} is the real one, read from the same maps {@code
* CallerResolver#resolve} consults, so a target that classifier calls known is one {@code
* resolve} would actually resolve as {@link Role#OBSERVER}.
*/
public static final Predicate<String> NO_KNOWN_OBSERVER_TARGET = target -> false;
/**
* Convenience form for a caller with no classifier to supply. Fails closed: a collaborator's
* {@code SEND} is refused, as if no terminal were a configured lead or collaborator — the
* same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR} gives explicitly. Every other action's
* result is identical to the four-argument form's, since none of them consult the classifier.
* or an observer's {@code SEND} is refused, as if no terminal were a configured lead,
* collaborator, or observer target — the same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR}
* and {@link #NO_KNOWN_OBSERVER_TARGET} give explicitly. Every other action's result is
* identical to the five-argument form's, since none of them consult either classifier.
*
* <p>Its default classifier denies every collaborator, so a caller enforcing authorization
* must use the four-argument form instead.
* <p>Its default classifiers deny every collaborator and every observer, so a caller
* enforcing authorization must use the five-argument form instead.
*/
public static boolean permits(Principal caller, Action action, String targetSession) {
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR);
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR, NO_KNOWN_OBSERVER_TARGET);
}
/**
* As {@link #permits(Principal, Action, String)}, with a real classifier for a collaborator's
* {@code SEND}. An observer's {@code SEND} still fails closed ({@link #NO_KNOWN_OBSERVER_TARGET}) —
* a caller enforcing both grants must use the five-argument form.
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator) {
return permits(caller, action, targetSession, knownLeadOrCollaborator, NO_KNOWN_OBSERVER_TARGET);
}
/**
* Whether {@code caller} may perform {@code action} against {@code targetSession}.
*
* @param targetSession the session id in the request path; only consulted for the
* worker-scoped actions ({@code REPLY}, {@code ASK}) and for a
* collaborator's {@code SEND}, ignored otherwise, may be
* {@code null}
* worker-scoped actions ({@code REPLY}, {@code ASK}), for a
* collaborator's {@code SEND}, and for an observer's
* {@code SEND}, ignored otherwise, may be {@code null}
* @param knownLeadOrCollaborator whether a terminal is a configured lead or collaborator —
* consulted only for a collaborator's {@code SEND}, to confine
* it to another named peer and never a spawned member's
* terminal
* @param knownObserverTarget whether a terminal is one this daemon would itself resolve as
* {@link Role#OBSERVER} — consulted only for an observer's
* {@code SEND}, to confine it to another observer pane and never
* a lead, a collaborator, or a spawned member
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator) {
Predicate<String> knownLeadOrCollaborator,
Predicate<String> knownObserverTarget) {
if (caller == null || caller.isAnonymous()) {
return false; // authenticated as nothing ⇒ authorized for nothing
}
@@ -107,13 +132,15 @@ public final class Authz {
// would be a worker escalating into the orchestrator role.
case SPAWN, STOP, DRAIN, HANDOVER -> caller.isPrimary();
// Delivering a turn to a local session is open to the primary, the architect, and a
// collaborator whose target is itself a configured lead or collaborator: the architect
// delegates to workers (that is the role's point); a collaborator may reach only
// another named peer, never a spawned member's terminal. A worker is excluded —
// sending would be it escalating.
// Delivering a turn to a local session is open to the primary and the architect
// unconditionally. A collaborator may reach only a target that is itself a configured
// lead or collaborator, never a spawned member's terminal. An observer may reach only
// a target that would itself resolve as an observer, never a lead, a collaborator, or
// a spawned member. A worker is excluded from every case — sending would be it
// escalating into the orchestrator role.
case SEND -> caller.isPrimary() || caller.isArchitect()
|| (caller.isCollaborator() && knownLeadOrCollaborator.test(targetSession));
|| (caller.isCollaborator() && knownLeadOrCollaborator.test(targetSession))
|| (caller.isObserver() && knownObserverTarget.test(targetSession));
// Resolving a worker's blocked question is part of delegating to it, open to the same
// two roles that may stand up that delegation in the first place. Not a collaborator:
@@ -270,6 +270,31 @@ public final class CallerResolver {
|| collaboratorTerminals.get().containsKey(target);
}
/**
* Whether {@code target} names a terminal this resolver would itself resolve as {@link
* Role#OBSERVER} — the classifier an observer's {@code SEND} is checked against, read from the
* same maps and functions {@link #resolve} consults so a target this accepts is exactly one
* {@code resolve} would hand back {@link Role#OBSERVER} for, and the reverse.
*/
public Predicate<String> sendableObserverTarget() {
return target -> target != null
&& spawnedMemberRole.apply(target) == null
&& !leadTerminals.get().containsKey(target)
&& !boundToArchitectSlot(target)
&& !collaboratorTerminals.get().containsKey(target);
}
/**
* Whether {@code terminal} is bound to a configured slot the live roster still confirms as an
* architect — the one classifier {@link #sendableObserverTarget()} and {@code FleetMcp}'s
* {@code panes} row both read, so a pane's reported role and its {@code SEND} reachability can
* never drift apart.
*/
public boolean boundToArchitectSlot(String terminal) {
String slot = architectTerminals.get().get(terminal);
return slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT;
}
/**
* Resolve the caller of a request.
*
@@ -146,8 +146,9 @@ public record Principal(Role role, String terminal, long pid, String name) {
}
/**
* 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.
* Stable identity used to own tickets and open turns. The unnamed primary has no owner key —
* {@code null} — and that is matched against a ticket's recorded owner the same way any other
* key is: it owns a ticket another unnamed primary created, and nothing else.
*/
public String ownerKey() {
return switch (role) {
@@ -12,9 +12,10 @@ package dev.ltms.fleet.auth;
public enum Role {
/**
* The orchestrating session. Established either by being a loopback caller that is not a
* worker pane (under {@code loopback-trust}) or by presenting a valid bearer token (under
* {@code token} mode).
* The orchestrating session. Established either by being a loopback caller that resolves to
* no herdr pane at all (under {@code loopback-trust}) or by presenting a valid bearer token
* (under {@code token} mode). A loopback caller that does own a pane, but matches none of the
* roles below, resolves to {@link #OBSERVER} instead.
*/
PRIMARY,
@@ -49,10 +50,11 @@ public enum Role {
* 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}).
* honoured regardless of auth mode. May {@code READ} and {@code METRICS}, {@code REPLY}/
* {@code ASK} only as its own pane, and {@code SEND} only to a target that would itself
* resolve as {@code OBSERVER}; may not {@code SPAWN}/{@code STOP}/{@code DRAIN}/
* {@code HANDOVER}, poll a ticket ({@code TASK_READ}), or reach the coordination broker
* ({@code COORD_SEND}/{@code COORD_READ}).
*/
OBSERVER,
@@ -28,6 +28,7 @@ import java.util.Comparator;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Set;
import java.util.regex.Pattern;
@@ -1092,8 +1093,8 @@ public record FleetConfig(
}
/**
* One entry of the CB-530 {@code leaders:} registry — a pane that orchestrates rather than one
* that is orchestrated.
* One entry of the {@code leaders:} registry — a pane that orchestrates rather than one that is
* orchestrated.
*
* <p>Why a registry and not a second {@code primary:}: {@code primary.terminal} is singular by
* construction, so a session in any other pane resolves as a worker. That is correct while one
@@ -1103,46 +1104,49 @@ public record FleetConfig(
* <p>{@code kind} and {@code model} are descriptive only: they document what runs in the pane
* and are reported back by {@code fleet_whoami}.
*
* <p><b>A lead is now also creatable (CB-557).</b> Before, nothing spawned one — a lead
* pre-existed, which is why it had to be recognised by configuration rather than created. With
* {@code profile} and {@code instances} the daemon may stand one up when none is live, so the
* pane no longer has to exist before the daemon does. Recognition still comes first: a lead
* already running in its configured {@code tab} is adopted, and only the shortfall is launched.
* <p>A lead with {@code profile} and {@code instances} set may be launched by the daemon when
* none is live; recognition always comes first, so only the shortfall is launched.
*
* <p><b>{@code tab} replaced {@code terminal} (CB-579).</b> A herdr {@code terminal_id} changes
* every time the lead's session restarts, so pinning one cost a config edit and a daemon restart
* per restart. A tab is stable: a human opens it once, it holds exactly one pane, and its label
* survives restarts of the agent inside it — so identity is now the tab label alone.
* <p>Every lead's tab is labelled {@link #LEAD_TAB_LABEL}, a fixed constant — not a per-entry
* config value. {@code workspace} is therefore what tells one lead from another: two leaders
* sharing one space would both resolve to the one tab named {@code lead} there, so only one
* could ever be found. {@code tab} is a deprecated legacy label, still matched within this
* lead's own space alongside the constant.
*
* @param profile the {@code profiles:} entry to launch this lead on when one must
* be created; {@code null} ⇒ recognise-only, never create.
* <p>fleetd #176: also the field {@code Fleetd.leadSeatLookup} reads
* to learn which account this lead's own live session shares — set it
* <p>Also the field {@code Fleetd.leadSeatLookup} reads to learn
* which account this lead's own live session shares — set it
* (safely, even on an already-running recognise-only lead: naming a
* profile here never starts anything beyond {@code instances}) so a
* {@code subscription: true} worker profile sharing its
* {@code effectiveCredentialId()} has this lead's seat subtracted from
* {@code fleet_list}'s {@code free}. {@code null} here also means this
* lead's seat cannot be derived and is not counted.
* @param tab the exact tab label hosting this lead, matched case-insensitively;
* the only field identity depends on. Required — a lead with no
* {@code tab} can never be discovered, launched or not
* @param tab deprecated legacy tab label, matched case-insensitively within
* this lead's own space alongside {@link #LEAD_TAB_LABEL}. Optional —
* {@code null}/blank means only the constant is accepted
* @param instances how many of this lead should be live (default 1). The daemon
* launches only the shortfall, so a restart adopts rather than doubles
* @param tabPrefix lead-tab naming convention checked against member labels. Lead
* identity uses {@code tab}. Default {@code "lead:"}
* identity uses {@link #acceptedLabels()}. Default {@code "lead:"}
* @param scanIntervalSeconds how long a tab scan is cached before herdr is asked again; also the
* worst case before a newly-labelled tab is recognised. Default 10
* @param kind which agent runs there ({@code claude}, {@code opencode}, …)
* @param model the model or selector it runs, for operators reading the roster
* @param workspace the space this lead's tab lives in — the uniqueness boundary
* identity now depends on. Default {@link #DEFAULT_WORKSPACE}
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Leader(String profile, String tab, Integer instances, String tabPrefix,
Integer scanIntervalSeconds, String kind, String model,
String workspace, String cwd) {
/** The tab label every lead is found by, and an auto-launched instance is created with. */
public static final String LEAD_TAB_LABEL = "lead";
/**
* Where an auto-launched lead's tab is created (CB-558). It defaults to the SAME shared
* Where an auto-launched lead's tab is created. It defaults to the SAME shared
* {@code "fleet"} space the members use, so the operator sees one "session" with many tabs.
* The scanner no longer excludes member spaces — it tells a lead from a member by the exact
* tab label, so a lead sharing the members' space is still discovered (see LeadLauncher).
@@ -1170,9 +1174,23 @@ public record FleetConfig(
return profile != null && !profile.isBlank() && instances > 0;
}
/** The tab label an auto-launched instance of this lead gets — its configured {@code tab}. */
/** The tab label an auto-launched instance of this lead gets. */
public String tabLabel() {
return tab;
return LEAD_TAB_LABEL;
}
/**
* The normalised labels (stripped, lower-cased) a tab in this lead's own space may carry to
* be recognised as this lead: {@link #LEAD_TAB_LABEL} first, plus the deprecated {@code tab}
* when configured and different. Every matcher in this class's callers reads this method —
* none re-derives the set.
*/
public List<String> acceptedLabels() {
String normalizedTab = (tab == null) ? null : tab.toLowerCase(Locale.ROOT);
if (normalizedTab == null || normalizedTab.equals(LEAD_TAB_LABEL)) {
return List.of(LEAD_TAB_LABEL);
}
return List.of(LEAD_TAB_LABEL, normalizedTab);
}
}
@@ -2095,13 +2113,6 @@ public record FleetConfig(
}
}
/**
* The top-level keys in {@code yaml} that this build does not understand, sorted. Package-private
* so the guardrail is asserted directly rather than through a log appender.
*
* @return empty when everything is known, or when {@code yaml} is not a mapping at all (a
* malformed file is {@code readValue}'s error to report, not this method's)
*/
/**
* Top-level keys renamed by the member taxonomy, mapped old → new.
*
@@ -2636,6 +2647,13 @@ public record FleetConfig(
}
}
/**
* The top-level keys in {@code yaml} that this build does not understand, sorted. Package-private
* so the guardrail is asserted directly rather than through a log appender.
*
* @return empty when everything is known, or when {@code yaml} is not a mapping at all (a
* malformed file is {@code readValue}'s error to report, not this method's)
*/
static List<String> unknownTopLevelKeys(String yaml) {
Map<?, ?> raw;
try {
@@ -2747,23 +2765,42 @@ public record FleetConfig(
/**
* Reject a member tab-label template that could render as a configured lead or collaborator
* tab or match a lead-tab naming convention, and reject two {@code fleet.leaders} or
* {@code fleet.collaborators} entries — across either registry — that share one exact tab.
* tab, as the fixed lead tab label, or that matches a lead-tab naming convention; reject two
* {@code fleet.leaders} entries that share one space; reject two {@code fleet.collaborators}
* entries — or a lead and a collaborator — that share one exact tab; and reject a collaborator
* tab equal to the fixed lead tab label.
*
* <p>{@code fleet.collaborators} has no {@code tabPrefix}: identity is matched on the exact
* {@code tab} alone, so only the exact-render check applies there, not the prefix check.
*
* @throws IllegalStateException when the fleet template or a profile {@code tabLabel} override
* can render as a configured lead or collaborator tab or match a
* lead-tab prefix, or when two entries — of either registry, or
* one of each — carry the same exact {@code tab}
* (case-insensitively)
* can render as a configured lead or collaborator tab, as the
* fixed lead tab label, or match a lead-tab prefix; when two
* leaders share one space; when two collaborators (or a lead and
* a collaborator) carry the same exact {@code tab}
* (case-insensitively); or when a collaborator's {@code tab}
* equals the fixed lead tab label
*/
public void validateLeadTabPrefixes() {
if (fleet == null) {
return;
}
List<String> bad = new ArrayList<>();
if (templateCanRenderAs(fleet.tabLabel(), Leader.LEAD_TAB_LABEL)) {
bad.add("fleet.tabLabel=\"" + fleet.tabLabel() + "\" can render as \""
+ Leader.LEAD_TAB_LABEL + "\", the fixed lead tab label");
}
profiles().entrySet().stream()
.map(Map.Entry::getKey)
.sorted()
.forEach(p -> {
String label = profiles().get(p).tabLabel();
if (templateCanRenderAs(label, Leader.LEAD_TAB_LABEL)) {
bad.add("profile '" + p + "' overrides tabLabel with \"" + label
+ "\", which can render as \"" + Leader.LEAD_TAB_LABEL
+ "\", the fixed lead tab label");
}
});
fleet.leaders().forEach((leadName, leader) -> {
if (leader == null) {
return;
@@ -2798,6 +2835,10 @@ public record FleetConfig(
return;
}
String tab = collaborator.tab();
if (tab != null && tab.equalsIgnoreCase(Leader.LEAD_TAB_LABEL)) {
bad.add("fleet.collaborators." + collabName + ".tab=\"" + tab + "\" is the fixed "
+ "lead tab label — a collaborator there would shadow a lead");
}
if (templateCanRenderAs(fleet.tabLabel(), tab)) {
bad.add("fleet.tabLabel=\"" + fleet.tabLabel() + "\" can render as the tab of "
+ "collaborator '" + collabName + "' (\"" + tab + "\")");
@@ -2827,18 +2868,20 @@ public record FleetConfig(
for (int i = 0; i < leadNames.size(); i++) {
String nameA = leadNames.get(i);
Leader a = fleet.leaders().get(nameA);
if (a == null || a.tab() == null || a.tab().isBlank()) {
if (a == null) {
continue;
}
for (int j = i + 1; j < leadNames.size(); j++) {
String nameB = leadNames.get(j);
Leader b = fleet.leaders().get(nameB);
if (b == null || b.tab() == null || b.tab().isBlank()) {
if (b == null) {
continue;
}
if (a.tab().equalsIgnoreCase(b.tab())) {
collisions.add("lead '" + nameA + "' and lead '" + nameB + "' both use tab \""
+ a.tab() + "\"");
if (a.workspace().equalsIgnoreCase(b.workspace())) {
collisions.add("lead '" + nameA + "' and lead '" + nameB + "' share workspace \""
+ a.workspace() + "\" — both would resolve to the tab named \""
+ Leader.LEAD_TAB_LABEL + "\" in that space, so only one could ever be "
+ "found");
}
}
}
@@ -2882,9 +2925,9 @@ public record FleetConfig(
return;
}
throw new IllegalStateException("refusing to start: " + String.join("; ", collisions)
+ ". Tab identity is matched exactly, so only one of two entries sharing a tab can "
+ "ever be found — the other is silently unreachable. Give each lead and "
+ "collaborator its own exact tab.");
+ ". Identity is matched exactly, so only one of two entries sharing a space or a "
+ "tab can ever be found — the other is silently unreachable. Give each lead its "
+ "own space, and each collaborator its own exact tab.");
}
private static boolean templateCanRenderAs(String template, String tab) {
@@ -2913,30 +2956,32 @@ public record FleetConfig(
}
/**
* Reject a profile that places its members by {@code "pane"} while any {@code fleet.leaders}
* or {@code fleet.collaborators} entry names a {@code tab}. A pane-placed member lands inside
* the focused tab rather than its own, so it can land inside a lead's or collaborator's own
* labelled tab. {@link dev.ltms.fleet.herdr.LeadTabScanner} identifies a lead or collaborator
* purely by that tab's label — it does not exclude the member space — so a member that ends up
* there, while its pane carries no entry in the spawned-member roster, is read back as that
* lead or collaborator and granted that identity's authority.
* Reject a profile that places its members by {@code "pane"} while {@code fleet.leaders} has
* any entry, or any {@code fleet.collaborators} entry names a {@code tab}. A pane-placed
* member lands inside the focused tab rather than its own, so it can land inside a lead's or
* collaborator's own labelled tab. {@link dev.ltms.fleet.herdr.LeadTabScanner} identifies a
* lead or collaborator purely by that tab's label — it does not exclude the member space — so
* a member that ends up there, while its pane carries no entry in the spawned-member roster,
* is read back as that lead or collaborator and granted that identity's authority.
*
* <p>Only an entry with a non-blank {@code tab} is in scope: one with no {@code tab} feeds
* <p>Every {@code fleet.leaders} entry is in scope regardless of its own {@code tab} field:
* {@link Leader#acceptedLabels()} always includes {@link Leader#LEAD_TAB_LABEL}. Only a
* collaborator with a non-blank {@code tab} is in scope: one with no {@code tab} feeds
* nothing into {@link dev.ltms.fleet.herdr.LeadTabScanner}, so it creates no hazard here.
*
* @throws IllegalStateException when any {@code profiles:} entry is pane-placed while any
* {@code fleet.leaders} or {@code fleet.collaborators} entry
* names a non-blank {@code tab}
* @throws IllegalStateException when any {@code profiles:} entry is pane-placed while
* {@code fleet.leaders} is non-empty, or any
* {@code fleet.collaborators} entry names a non-blank
* {@code tab}
*/
public void validatePanePlacementAgainstLeadTabs() {
if (fleet == null) {
return;
}
boolean anyLeaderHasTab = fleet.leaders().values().stream()
.anyMatch(leader -> leader != null && leader.tab() != null && !leader.tab().isBlank());
boolean anyLead = !fleet.leaders().isEmpty();
boolean anyCollaboratorHasTab = fleet.collaborators().values().stream()
.anyMatch(c -> c != null && c.tab() != null && !c.tab().isBlank());
if (!anyLeaderHasTab && !anyCollaboratorHasTab) {
if (!anyLead && !anyCollaboratorHasTab) {
return;
}
List<String> bad = new ArrayList<>();
@@ -2953,8 +2998,9 @@ public record FleetConfig(
+ "pane-placed member can land inside that labelled tab, and while its pane "
+ "carries no entry in the spawned-member roster, it is read back as the lead or "
+ "collaborator and granted that identity's authority. Set placement: tab for "
+ "each named profile, or remove the tab from every fleet.leaders and "
+ "fleet.collaborators entry.");
+ "each named profile — the only fix when a lead triggered this, since a lead's "
+ "tab label is fixed regardless of its own tab: field. A collaborator's tab can "
+ "still be removed instead.");
}
/**
@@ -3061,14 +3107,13 @@ public record FleetConfig(
* so duplicates are unrepresentable by construction once loaded — and {@link #load(Path)}
* already rejects a duplicated slot name at parse time, before the map collapses.
*
* <p>Also rejects a {@code fleet.collaborators} entry with no (or a blank) {@code tab}. A
* {@code profile}-less lead is still useful recognise-only — {@code tab} is the only field
* that matters to it either way. A collaborator carries no other field at all, so a blank
* {@code tab} leaves nothing for the entry to mean.
* <p>A lead's {@code profile} is optional — a {@code profile}-less lead is still useful
* recognise-only. Also rejects a {@code fleet.collaborators} entry with no (or a blank)
* {@code tab}: a collaborator carries no other field at all, so a blank {@code tab} leaves
* nothing for the entry to mean.
*
* @throws IllegalStateException when a slot names no profile or an unknown one, when a lead
* can be neither found nor created, or when a collaborator names
* no tab, naming the offending entry
* @throws IllegalStateException when a slot or a lead references an unknown profile, or a
* collaborator names no tab, naming the offending entry
*/
public void validateMembers() {
if (fleet == null) {
@@ -3101,10 +3146,6 @@ public record FleetConfig(
+ "', which is not a configured profiles: entry (have: " + profiles.keySet()
+ ").");
}
if (leader.tab() == null || leader.tab().isBlank()) {
bad.add("fleet.leaders." + name + " has no tab: — a lead is now found (and, if "
+ "auto-launched, labelled) purely by its tab, so every entry must name one.");
}
});
fleet.collaborators().forEach((name, collaborator) -> {
if (collaborator == null) {
@@ -25,14 +25,12 @@ import java.util.function.Supplier;
* by first starting the session and asking it. Scanning closes that loop: label the tab, and the
* pane is recognised on the next resolve.
*
* <p><strong>CB-579 — matched by name, not prefix.</strong> This used to strip one shared
* {@code tabPrefix} off a label to derive the lead's name, and merged a config-supplied
* {@code terminal_id} pin over every scan result so the pin could never expire. Both are gone: each
* lead now configures its own exact {@code tab} label ({@code fleet.leaders.<name>.tab}), so this
* class is handed a {@code tab → name} map up front and matches labels against it exactly
* (case-insensitively). There is no merge step — a scan result is the whole answer. That is the
* fix for the bug this replaces: a {@code terminal_id} pin surviving in config after the pane it
* named was gone, so the daemon kept treating a dead session as a live lead forever.
* <p><strong>Matched by label within a space, not by a shared prefix.</strong> Each lead's accepted
* labels (the fixed {@code lead} label, plus a deprecated {@code tab} when still configured) are
* matched exactly (case-insensitively) against tabs in that lead's own space only — a tab named
* {@code lead} in one space never resolves to another space's lead. A scan result is the whole
* answer; nothing is merged in from configuration between scans, so a tab that is gone drops out on
* the very next scan instead of lingering forever.
*
* <p><strong>Direction of trust.</strong> The label names the lead; it never <em>grants</em>
* anything a pane could take for itself. Three properties keep that honest:
@@ -103,7 +101,8 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
private record Entry(String name, Kind kind) {}
private final HerdrClient herdr;
private final Map<String, Entry> tabToEntry;
private final Map<String, Map<String, String>> leadLabelsBySpace;
private final Map<String, String> collaboratorTabToName;
private final Set<String> excludedWorkspaceLabels;
private final long ttlNanos;
private final LongSupplier clock;
@@ -124,15 +123,17 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
/**
* @param herdr the herdr client to query ({@code workspace.list},
* {@code tab.list}, {@code pane.list} — all read-only)
* @param tabToName every configured lead's exact tab label → its name
* ({@code fleet.leaders.<name>.tab}), matched case-insensitively
* @param leadLabelsBySpace each configured lead's accepted tab labels, keyed by the
* lead's own space label, then by label, to its name — matched
* case-insensitively on both the space and the label. A tab
* matches a lead only within that lead's own space
* @param excludedWorkspaceLabels workspaces never scanned — the configured worker spaces
* @param ttlNanos how long a scan result is reused before the next one
* @param clock nanosecond time source ({@code System::nanoTime} in production)
*/
public LeadTabScanner(HerdrClient herdr, Map<String, String> tabToName,
public LeadTabScanner(HerdrClient herdr, Map<String, Map<String, String>> leadLabelsBySpace,
Set<String> excludedWorkspaceLabels, long ttlNanos, LongSupplier clock) {
this(herdr, tabToName, Map.of(), excludedWorkspaceLabels, ttlNanos, clock);
this(herdr, leadLabelsBySpace, Map.of(), excludedWorkspaceLabels, ttlNanos, clock);
}
/**
@@ -140,14 +141,15 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
* for configured collaborator tabs in the same pass.
*
* @param collaboratorTabToName every configured collaborator's exact tab label → its name
* ({@code fleet.collaborators.<name>.tab}), matched the same way as
* {@code tabToName}
* ({@code fleet.collaborators.<name>.tab}), matched
* case-insensitively in any space
*/
public LeadTabScanner(HerdrClient herdr, Map<String, String> tabToName,
public LeadTabScanner(HerdrClient herdr, Map<String, Map<String, String>> leadLabelsBySpace,
Map<String, String> collaboratorTabToName,
Set<String> excludedWorkspaceLabels, long ttlNanos, LongSupplier clock) {
this.herdr = herdr;
this.tabToEntry = buildTabIndex(tabToName, collaboratorTabToName);
this.leadLabelsBySpace = buildLeadIndex(leadLabelsBySpace);
this.collaboratorTabToName = normalizedLabelMap(collaboratorTabToName);
this.excludedWorkspaceLabels = excludedWorkspaceLabels == null
? Set.of() : Set.copyOf(excludedWorkspaceLabels);
this.ttlNanos = ttlNanos;
@@ -155,30 +157,46 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
}
/**
* Keys stripped and lower-cased once, so every lookup is a plain map hit. Leads and
* collaborators merge into a single index, so {@link #scan()} matches both kinds in one pass
* over the tab list; a label naming both a lead and a collaborator takes the lead entry —
* leads are put last, so a colliding key's lead entry is the one that overwrites — since a lead
* can already do everything a collaborator can. Config validation already refuses a lead and a
* collaborator sharing one exact tab, so this ordering is defence in depth, not the control.
* Space and label keys stripped and lower-cased once, so every lookup is a plain map hit. A
* space with no usable labels is simply absent — {@link #leadLabelsFor} then finds nothing for
* it, which is also what a space with a {@code null} label gets.
*/
private static Map<String, Entry> buildTabIndex(Map<String, String> tabToName,
Map<String, String> collaboratorTabToName) {
Map<String, Entry> out = new LinkedHashMap<>();
putNormalized(out, collaboratorTabToName, Kind.COLLABORATOR);
putNormalized(out, tabToName, Kind.LEAD);
private static Map<String, Map<String, String>> buildLeadIndex(
Map<String, Map<String, String>> leadLabelsBySpace) {
Map<String, Map<String, String>> out = new LinkedHashMap<>();
if (leadLabelsBySpace == null) {
return Map.of();
}
leadLabelsBySpace.forEach((space, labelsToName) -> {
if (space == null || space.isBlank()) {
return;
}
Map<String, String> normalized = normalizedLabelMap(labelsToName);
if (!normalized.isEmpty()) {
out.put(space.strip().toLowerCase(Locale.ROOT), normalized);
}
});
return Collections.unmodifiableMap(out);
}
private static void putNormalized(Map<String, Entry> out, Map<String, String> tabToName, Kind kind) {
if (tabToName == null) {
return;
private static Map<String, String> normalizedLabelMap(Map<String, String> labelToName) {
Map<String, String> out = new LinkedHashMap<>();
if (labelToName != null) {
labelToName.forEach((label, name) -> {
if (label != null && !label.isBlank() && name != null && !name.isBlank()) {
out.put(label.strip().toLowerCase(Locale.ROOT), name);
}
});
}
tabToName.forEach((tab, name) -> {
if (tab != null && !tab.isBlank() && name != null && !name.isBlank()) {
out.put(tab.strip().toLowerCase(Locale.ROOT), new Entry(name, kind));
}
});
return Collections.unmodifiableMap(out);
}
/** The accepted lead labels configured for {@code spaceLabel}, or an empty map for no match. */
private Map<String, String> leadLabelsFor(String spaceLabel) {
if (spaceLabel == null) {
return Map.of();
}
return leadLabelsBySpace.getOrDefault(spaceLabel.strip().toLowerCase(Locale.ROOT), Map.of());
}
/**
@@ -241,9 +259,10 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
if (ws.workspaceId() == null || excludedWorkspaceLabels.contains(ws.label())) {
continue;
}
Map<String, String> leadLabelsHere = leadLabelsFor(ws.label());
for (JsonNode t : herdr.call("tab.list", Map.of("workspace_id", ws.workspaceId())).path("tabs")) {
Tab tab = Tab.from(t);
Entry entry = entryOf(tab.label());
Entry entry = entryOf(tab.label(), leadLabelsHere);
if (entry != null && tab.tabId() != null) {
entryByTab.put(tab.tabId(), entry);
}
@@ -296,19 +315,27 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
}
/**
* The entry a tab label declares, or {@code null} if it names neither a configured lead nor a
* configured collaborator.
* The entry a tab label declares within one space, or {@code null} if it names neither a lead
* accepted in {@code leadLabelsHere} nor a configured collaborator.
*
* <p>Exact match (case-insensitive, ends stripped) against {@link #tabToEntry} — no prefix
* stripping, so an operator's {@code "lead: something-else"} tab is never mistaken for a
* configured lead just because it shares a prefix. The match strips a trailing
* {@link PendingCloseMarker} first, so a tab {@code LeadLauncher} has flagged as maybe-dead but
* not yet closed keeps resolving normally while that reconcile is pending.
* <p>Exact match (case-insensitive, ends stripped) — no prefix stripping, so an operator's
* {@code "lead: something-else"} tab is never mistaken for a configured lead just because it
* shares a prefix. The match strips a trailing {@link PendingCloseMarker} first, so a tab
* {@code LeadLauncher} has flagged as maybe-dead but not yet closed keeps resolving normally
* while that reconcile is pending. A lead match wins over a collaborator match for the same
* label — a lead can already do everything a collaborator can, and config validation refuses a
* lead and a collaborator sharing one exact tab in the first place.
*/
private Entry entryOf(String label) {
private Entry entryOf(String label, Map<String, String> leadLabelsHere) {
if (label == null) {
return null;
}
return tabToEntry.get(PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT));
String normalized = PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT);
String leadName = leadLabelsHere.get(normalized);
if (leadName != null) {
return new Entry(leadName, Kind.LEAD);
}
String collaboratorName = collaboratorTabToName.get(normalized);
return collaboratorName == null ? null : new Entry(collaboratorName, Kind.COLLABORATOR);
}
}
@@ -4,6 +4,7 @@ import com.fasterxml.jackson.databind.JsonNode;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
@@ -145,6 +146,52 @@ public final class PaneLocator {
return ancestry;
}
/**
* Every tab herdr tracks across every searched daemon, keyed by tab id, to its display label —
* the pane-discovery surface behind {@code GET /agents} and {@code fleet_list}'s {@code panes}
* row. Collapses to one scan in the single-daemon deployment, the same as
* {@link #terminalForPid}. A tab herdr reports with no label maps to a {@code null} value here;
* a tab with no {@code tab_id} is skipped.
*/
public Map<String, String> tabLabelsByTabId() {
Map<String, String> out = new LinkedHashMap<>();
for (HerdrClient herdr : herdrs) {
for (JsonNode w : herdr.call("workspace.list").path("workspaces")) {
String workspaceId = w.path("workspace_id").asText(null);
if (workspaceId == null) {
continue;
}
for (JsonNode t : herdr.call("tab.list", Map.of("workspace_id", workspaceId)).path("tabs")) {
Tab tab = Tab.from(t);
if (tab.tabId() != null) {
out.put(tab.tabId(), tab.label());
}
}
}
}
return out;
}
/**
* Every workspace ("space") herdr tracks across every searched daemon, keyed by workspace id,
* to its display label — the human-readable name behind {@code fleet_list}'s {@code panes} row,
* next to herdr's own internal {@code workspaceId}. Collapses to one scan in the single-daemon
* deployment, the same as {@link #terminalForPid}. A workspace herdr reports with no label maps
* to a {@code null} value here; a workspace with no {@code workspace_id} is skipped.
*/
public Map<String, String> workspaceLabelsByWorkspaceId() {
Map<String, String> out = new LinkedHashMap<>();
for (HerdrClient herdr : herdrs) {
for (JsonNode w : herdr.call("workspace.list").path("workspaces")) {
Workspace workspace = Workspace.from(w);
if (workspace.workspaceId() != null) {
out.put(workspace.workspaceId(), workspace.label());
}
}
}
return out;
}
/** Whether a pane owns one of the scanned pid's ancestors, or the check of it failed outright. */
private enum Ownership { OWNS, DOES_NOT_OWN, UNKNOWN }
@@ -0,0 +1,169 @@
package dev.ltms.fleet.herdr;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* Whether an agent pane's input box is clear for a delivery.
*
* <p>{@link AgentControl#send} pastes its text and submits it in the same call, so a delivery into a
* pane whose input box already holds characters submits those characters too. {@link
* AgentStatus#injectable()} cannot see that: it describes the agent, and an agent waiting at its
* prompt reports the same status whether its box is empty or holds a half-typed line. This reads the
* box itself.
*
* <p>Only a box that is positively empty clears the gate. A box with content, a pane this cannot
* recognise, and a failed read all hold the delivery, because a held delivery is recoverable and a
* submitted half-line is not. Every caller must therefore be a path that retries.
*
* <p>A pane that holds for {@link #HOLD_WARN_STREAK} consecutive checks gets one warning, so a box
* that never clears is visible instead of silent. The warning repeats only after the box has cleared
* again.
*/
public final class PromptBox {
private static final Logger log = LoggerFactory.getLogger(PromptBox.class);
/**
* herdr {@code agent.read} source. {@code detection} is the region herdr itself uses for status
* detection, so the input box is always drawn in it. It is not a short tail: it carries transcript
* scrollback above the box, including earlier prompts the operator has already submitted, which is
* why only the last box line on it is the live one.
*/
static final String PROBE_SOURCE = "detection";
/** Consecutive holds for one target before one warning is logged. */
static final int HOLD_WARN_STREAK = 20;
/**
* Input box markers, each matched only as a line's first characters: the caret the current TUI
* draws, and the bordered box an older one drew. A marker further along a line is transcript text,
* such as a caret inside something the operator quoted.
*/
private static final List<String> BOX_MARKERS = List.of("❯", "│ >");
/** Marker of a turn that is still generating; a box drawn under it is not a settled prompt. */
private static final String ACTIVE_TURN_MARKER = "esc to interrupt";
/** Block glyphs a terminal capture can leave in an otherwise empty box for the cursor cell. */
private static final String CURSOR_GLYPHS = "█▉▊▋▌▍▎▏";
/** What a box holds: nothing, unsubmitted characters, or a pane this cannot read as a box. */
public enum State { EMPTY, DRAFT, UNREADABLE }
/** A box reading: its state, and how many characters it holds ({@code 0} unless {@code DRAFT}). */
public record Reading(State state, int characters) {
}
private final AgentControl agents;
/** Consecutive holds per target, so a box that never clears can be warned about once. */
private final Map<String, Integer> holdStreaks = new ConcurrentHashMap<>();
public PromptBox(AgentControl agents) {
this.agents = agents;
}
/**
* Whether {@code target}'s input box is empty, so a delivery would submit only its own text.
* {@code false} means hold and come back; it never means the delivery failed.
*/
public boolean clearToSubmit(String target) {
Reading reading = inspect(target);
if (reading.state() == State.EMPTY) {
holdStreaks.remove(target);
return true;
}
int streak = holdStreaks.merge(target, 1, Integer::sum);
if (streak == HOLD_WARN_STREAK) {
log.warn("prompt box of {} has held a delivery {} times in a row ({}, {} character(s) in the box)"
+ " — nothing is lost, delivery resumes once the box is empty",
target, streak, reading.state(), reading.characters());
} else {
log.debug("prompt box of {} is {} ({} character(s)), holding delivery {}",
target, reading.state(), reading.characters(), streak);
}
return false;
}
/** Read and classify {@code target}'s pane. A read failure reads as {@link State#UNREADABLE}. */
private Reading inspect(String target) {
String pane;
try {
pane = agents.read(target, PROBE_SOURCE);
} catch (RuntimeException e) {
log.debug("prompt box read for {} failed, holding delivery: {}", target, e.getMessage());
return new Reading(State.UNREADABLE, 0);
}
return classify(pane);
}
/**
* Classify a Claude Code TUI pane region. Pure, so it is unit-testable without herdr.
*
* <p>{@link State#EMPTY} needs two positive signals: the pane's last box line holds nothing after
* its marker, and nothing below that line says a turn is still generating. Everything else is
* {@link State#UNREADABLE} — a blank capture, or a region with no box line at all — so a pane this
* does not understand holds the delivery rather than guessing it is safe.
*
* <p>The generating marker is looked for only from the box line down. Above it is scrollback, where
* an earlier turn's marker survives; treating that as a live turn would make {@link State#EMPTY}
* unreachable and hold every delivery forever.
*
* <p>Whitespace, a trailing box border and a cursor block count as nothing. Any other character
* counts as the operator's unsubmitted text, including a placeholder hint a future TUI might draw
* there; that direction holds a delivery it could have sent, which {@link #HOLD_WARN_STREAK} makes
* visible.
*/
static Reading classify(String pane) {
if (pane == null || pane.isBlank()) return new Reading(State.UNREADABLE, 0);
int box = lastBoxLineStart(pane);
if (box < 0) return new Reading(State.UNREADABLE, 0);
String fromBox = pane.substring(box);
if (fromBox.toLowerCase().contains(ACTIVE_TURN_MARKER)) return new Reading(State.UNREADABLE, 0);
String content = boxContent(firstLine(fromBox));
return content.isEmpty() ? new Reading(State.EMPTY, 0) : new Reading(State.DRAFT, content.length());
}
/** Offset of the last line starting with a box marker, or {@code -1} if the region has none. */
private static int lastBoxLineStart(String pane) {
int found = -1;
for (int start = 0; start <= pane.length(); ) {
int end = pane.indexOf('\n', start);
String line = pane.substring(start, end < 0 ? pane.length() : end);
if (markerLength(line) > 0) found = start;
if (end < 0) break;
start = end + 1;
}
return found;
}
/** Length of the box marker this line starts with, or {@code 0} if it starts with none. */
private static int markerLength(String line) {
for (String marker : BOX_MARKERS) {
if (line.startsWith(marker)) return marker.length();
}
return 0;
}
private static String firstLine(String text) {
int newline = text.indexOf('\n');
return newline < 0 ? text : text.substring(0, newline);
}
/** The text the box holds: its own line after the marker, stripped of border, padding and cursor. */
private static String boxContent(String boxLine) {
String line = boxLine.substring(markerLength(boxLine)).stripTrailing();
if (line.endsWith("│")) line = line.substring(0, line.length() - 1);
StringBuilder content = new StringBuilder();
for (char c : line.toCharArray()) {
if (Character.isWhitespace(c) || CURSOR_GLYPHS.indexOf(c) >= 0) continue;
content.append(c);
}
return content.toString();
}
}
@@ -4,16 +4,17 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.Set;
/**
* Tracks which workers are <em>available</em> — their Claude has booted and connected its MCP client
* to the bridge (CB-113). This is the reliable readiness signal, unlike herdr's {@code agent_status},
* which reports {@code idle} for a worker whose Claude is still booting. Delivering into that boot
* window pastes into a not-yet-ready TUI (the text is lost) and wedges the worker's delivery state,
* so the {@link Injector} holds the first delivery until the worker is present here.
* Tracks which peers are <em>available</em> — their Claude has booted and connected its MCP
* client to the bridge. For a spawned member this is the reliable readiness signal,
* unlike herdr's {@code agent_status}, which reports {@code idle} while its Claude is still
* booting. Delivering into that boot window pastes into a not-yet-ready TUI (the text is lost)
* and wedges that member's delivery state, so the {@link Injector} holds a spawned member's
* first delivery until it is present here.
*
* <p>Populated from the MCP transport: any MCP request whose connection resolves to a worker terminal
* marks that worker present (its {@code initialize} is the first such contact). A worker that never
* mounts the bridge MCP is never marked present — its sends stay queued until they time out, which is
* correct (it could not have replied anyway).
* <p>Populated from the MCP transport, for the peers whose deliverability rests on proving a live
* MCP contact rather than on a configured registry entry. A peer that never mounts the bridge MCP
* is never marked present — its sends stay queued until they time out, which is correct (it could
* not have replied anyway).
*/
public class MemberPresence {
@@ -18,6 +18,7 @@ import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
@@ -297,15 +298,11 @@ public final class LeadLauncher {
/**
* 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}
* (CB-579). A member sitting in the same shared workspace is not counted as a lead because its
* tab carries a different label, not because any workspace is excluded from this count.
*
* <p>There used to be a second path here — a running agent on the terminal a
* {@code fleet.leaders.<name>.terminal} pin named, for a lead opened and pinned by hand. That
* pin is retired: {@code tab} is now the only field identity depends on, and {@link Agent}
* already carries {@link Agent#tabId()} directly, so a hand-opened lead is found the same way an
* auto-launched one is — by labelling its tab to match.
* <em>not</em> live: a running agent in a tab, in that lead's own space, carrying one of its
* {@link FleetConfig.Leader#acceptedLabels()}. A member sitting in the same shared workspace is
* not counted as a lead because its tab carries a different label, not because any workspace is
* excluded from this count. A tab matching a lead's label in a <em>different</em> space is not
* counted either — space is the uniqueness boundary between leads.
*
* <p>fleetd #359 review finding 1: a labelled tab with nothing running in it is split into
* {@code toClose} (already flagged pending-close by a previous reconcile, and still dead — two
@@ -319,12 +316,11 @@ public final class LeadLauncher {
}
private Map<String, LeadCount> countLeads(Map<String, FleetConfig.Leader> leaders) {
// A lead and the members share ONE workspace now (the operator asked for a single "session"
// with many tabs), so a workspace can no longer be excluded wholesale — the lead lives in the
// member workspace by design. The sole discriminator is the exact tab label: a lead carries
// its configured `fleet.leaders.<name>.tab` ("lead: opus"), while a member carries its
// profile's `worker: {profile} #{n}` template. These never collide, so an exact-label match
// separates them without needing to know which workspace anyone is in.
// A lead and the members share ONE workspace (the operator asked for a single "session" with
// many tabs), so a workspace can no longer be excluded wholesale — the lead lives in the
// member workspace by design. The discriminator is the tab label together with the space: a
// member's tab never carries one of a lead's accepted labels, and a lead's own label only
// counts within that lead's configured space.
Map<String, String> nameByTab = new LinkedHashMap<>();
Set<String> flaggedTabIds = new LinkedHashSet<>();
for (Workspace ws : spaces.listWorkspaces()) {
@@ -332,7 +328,7 @@ public final class LeadLauncher {
continue;
}
for (Tab tab : spaces.listTabs(ws.workspaceId())) {
String declared = leadNameOf(tab.label(), leaders);
String declared = leadNameOf(tab.label(), ws.label(), leaders);
if (declared != null && tab.tabId() != null) {
nameByTab.put(tab.tabId(), declared);
if (PendingCloseMarker.isFlagged(tab.label())) {
@@ -382,21 +378,24 @@ public final class LeadLauncher {
}
/**
* The configured lead a tab label names, or {@code null} for a label that names none.
* The configured lead a tab names, or {@code null} for a label or space that names none.
*
* <p>Matched exactly (case-insensitively) against each lead's configured {@code tab}, so an
* operator's {@code "lead: something-else"} tab is not mistaken for a configured lead. A
* trailing {@link PendingCloseMarker} is stripped first, so a tab this class flagged on a
* previous reconcile is still recognised as the same lead's tab on this one.
* <p>A match requires both: the label (case-insensitively, trailing {@link PendingCloseMarker}
* stripped) must be one of the lead's {@link FleetConfig.Leader#acceptedLabels()}, and {@code
* space} must be that lead's own {@link FleetConfig.Leader#workspace()}. The same label in a
* different space names no lead — space is the uniqueness boundary between leads.
*/
private String leadNameOf(String label, Map<String, FleetConfig.Leader> leaders) {
if (label == null) {
private String leadNameOf(String label, String space, Map<String, FleetConfig.Leader> leaders) {
if (label == null || space == null) {
return null;
}
String l = PendingCloseMarker.strip(label);
String l = PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT);
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
String tab = e.getValue().tabLabel();
if (tab != null && l.equalsIgnoreCase(tab.strip())) {
FleetConfig.Leader lead = e.getValue();
if (lead == null || !lead.workspace().equalsIgnoreCase(space)) {
continue;
}
if (lead.acceptedLabels().contains(l)) {
return e.getKey();
}
}
@@ -93,7 +93,7 @@ import java.util.function.Supplier;
*
* <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
* PrimaryRegistry#primaryTerminal()}. That is correct for a background loop with no caller (see
* PrimaryRegistry#currentPrimaryTerminal()}. That is correct for a background loop with no caller (see
* {@code dev.ltms.fleet.msg.LeadHeartbeatLoop}), but wrong here and a violation of this project's
* own charter invariant 3 — "identity comes from the connection, never an argument." This daemon
* can hold more than one labelled lead tab (see {@code LeadLauncher}'s fleetd #359 two-reading
@@ -149,9 +149,17 @@ public final class LeadRollover {
* calling lead's workspace, before storing it here; see this class's
* javadoc. This is the value the MCP layer hands back to the lead as
* "write your file here", so callers may rely on it always being absolute.
* @param rolloverKey the key {@link #confirm}'s single-flight claim is taken under. {@link
* #open} resolves this once, here, from {@code leadTerminal} while the
* calling lead is certainly still live: the lead's configured name, or
* {@code leadTerminal} itself when no name resolves. Carried rather than
* recomputed at release time, because by the time a roll's continuation
* releases its claim the OLD terminal may no longer resolve to any name at
* all — recomputing there would release a different key than the one the
* claim was taken under.
*/
public record PendingRollover(String token, String leadTerminal, String handoverPath,
long requestedAtMillis) {}
long requestedAtMillis, String rolloverKey) {}
/** Which check refused a {@link #confirm} call, named so a caller can act on it. */
public enum RefusalReason {
@@ -176,9 +184,9 @@ public final class LeadRollover {
*/
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.
* This lead already has a roll running: an earlier {@link #confirm} call claimed its
* single-flight key (see {@link PendingRollover#rolloverKey}) and that roll's continuation
* has not released it yet. {@code detail} names the key and the token that holds the claim.
*/
ROLL_ALREADY_RUNNING
}
@@ -340,9 +348,12 @@ public final class LeadRollover {
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}.
* the terminal names no currently-recognised lead. {@link #open} calls this on the calling
* lead's own terminal, while it is certainly still live, to resolve {@link
* PendingRollover#rolloverKey}. The deferred continuation also calls this, on the OLD terminal,
* before tearing it down, so it knows which lead to pass to {@link LeadLauncher#relaunch} — by
* that point the live roster may no longer contain the old terminal, so this lookup can return
* {@code null} here even though {@link #open}'s earlier call against the same terminal did not.
*/
private final Function<String, String> leadNameForTerminal;
/**
@@ -362,14 +373,14 @@ 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
* {@link PendingRollover#rolloverKey} → the token of the roll currently holding that key
* 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.
* releases it in a {@code finally}, on both the success and the thrown-exception path. A key
* absent from this map has no roll currently in flight for it.
*/
private final Map<String, String> rollingByTerminal = new ConcurrentHashMap<>();
private final Map<String, String> rollingByLead = 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
@@ -438,8 +449,9 @@ public final class LeadRollover {
/**
* The lead says it is ready to be replaced. Generates a token and records the resolved
* handover path, this moment's wall-clock timestamp (the baseline {@link #confirm} checks the
* handover file's modified time against), and {@code leadTerminal} — only that exact terminal
* may later {@link #confirm} this token.
* handover file's modified time against), {@code leadTerminal} — only that exact terminal may
* later {@link #confirm} this token — and {@link PendingRollover#rolloverKey}, resolved here
* from {@code leadTerminal} while the calling lead is certainly still live.
*
* @param leadTerminal the calling lead's terminal id, resolved by the MCP layer from the
* connection (see this class's javadoc) — never a client-supplied value
@@ -459,7 +471,9 @@ public final class LeadRollover {
String token = UUID.randomUUID().toString();
long requestedAt = nowMillis.getAsLong();
String resolvedPath = resolveHandoverPath(cfg.handoverPath(), leadTerminal);
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt);
String leadName = leadNameForTerminal.apply(leadTerminal);
String rolloverKey = (leadName == null || leadName.isBlank()) ? leadTerminal : leadName;
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt, rolloverKey);
pending.put(token, p);
if (resolvedPath.equals(cfg.handoverPath())) {
log.info("lead-rollover: open token={} lead={} handoverPath={} reason={}",
@@ -565,12 +579,11 @@ public final class LeadRollover {
// 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);
// different, still-running roll already holds this lead's claim.
String holder = rollingByLead.putIfAbsent(p.rolloverKey(), token);
if (holder != null) {
return RollDecision.refused(RefusalReason.ROLL_ALREADY_RUNNING,
"lead terminal " + p.leadTerminal() + " already has a roll running under token "
+ holder);
"lead '" + p.rolloverKey() + "' already has a roll running under token " + holder);
}
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
@@ -591,14 +604,14 @@ public final class LeadRollover {
} 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.
// finally — the only other place that releases rollingByLead — 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.
// or this lead 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);
rollingByLead.remove(p.rolloverKey(), 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"));
@@ -641,10 +654,10 @@ public final class LeadRollover {
+ "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());
// path above — a release only on success would leave this lead 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 key.
rollingByLead.remove(p.rolloverKey(), p.token());
}
}
@@ -5,13 +5,13 @@ import dev.ltms.fleet.herdr.PaneLocator;
/**
* Resolves <em>who is calling</em> an MCP tool from the connection alone — the anti-spoofing
* identity model of the MCP contract. It ties the connection's loopback peer PID (from the OS)
* to a herdr agent pane (from herdr), yielding the caller's worker {@code terminal_id}. A caller
* that maps to no worker pane — the primary, or an off-host client — resolves to {@code null}.
* to a herdr agent pane (from herdr), yielding that pane's {@code terminal_id}. A connection that
* maps to no pane resolves to {@code null}; this class assigns no role to either outcome — {@link
* dev.ltms.fleet.auth.CallerResolver} does that.
*
* <p>Both sources are authoritative and unforgeable: the OS reports the real connecting PID, and
* herdr owns the PID→pane mapping. A worker cannot claim to be another worker, nor the primary.
* Single-host only (the herd shares the {@code fleetd} host); the token path is the split-host
* fallback.
* herdr owns the PID→pane mapping, so a caller cannot claim to be at another pane. Single-host
* only (the herd shares the {@code fleetd} host); the token path is the split-host fallback.
*/
public final class ConnectionIdentity {
@@ -46,9 +46,10 @@ public final class ConnectionIdentity {
}
/**
* The caller resolved from the connection: its worker {@code terminal} (or {@code null} for the
* primary / an off-host client), its {@code pid} (or {@code -1} if not resolvable), and whether
* the pane scan behind {@code terminal} ran to completion ({@link #scanComplete}).
* The caller resolved from the connection: the {@code terminal} of the pane it connects from
* (or {@code null} when the connection maps to no pane), its {@code pid} (or {@code -1} if not
* resolvable), and whether the pane scan behind {@code terminal} ran to completion
* ({@link #scanComplete}).
*/
public record Caller(String terminal, long pid, boolean scanComplete) {
@@ -87,8 +88,8 @@ public final class ConnectionIdentity {
}
/**
* The calling worker's {@code terminal_id}, or {@code null} if the caller is not a known
* on-host worker (treat as the primary).
* The terminal id of the pane the caller connects from, or {@code null} if the connection
* maps to no pane.
*/
public String callerTerminal(String remoteAddr, int remotePort) {
return resolve(remoteAddr, remotePort).terminal();
@@ -1,5 +1,6 @@
package dev.ltms.fleet.mcp;
import dev.ltms.fleet.Fleetd;
import dev.ltms.fleet.auth.AuditLog;
import dev.ltms.fleet.auth.Authz;
import dev.ltms.fleet.auth.CallerResolver;
@@ -41,6 +42,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import jakarta.servlet.http.HttpServlet;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -52,6 +54,7 @@ import java.util.concurrent.TimeUnit;
import java.util.function.BiFunction;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Predicate;
import java.util.function.Supplier;
import java.util.stream.Collectors;
@@ -116,9 +119,10 @@ public final class FleetMcp {
private final ConnectionIdentity identity;
/**
* Kept as a field (rather than only captured by the {@code contextExtractor} closure) so
* {@link #denyFor} can read {@link CallerResolver#knownLeadOrCollaborator()} — the classifier a
* collaborator's {@code SEND} is checked against, built from the same lead and collaborator
* maps {@link #identity}-based resolution reads.
* {@link #denyFor} can read {@link CallerResolver#knownLeadOrCollaborator()} and {@link
* CallerResolver#sendableObserverTarget()} — the classifiers a collaborator's and an
* observer's {@code SEND} are each checked against, built from the same maps {@link #identity}-
* based resolution reads.
*/
private final CallerResolver callers;
private final Metrics metrics; // CB-502: null → auth failures not counted
@@ -302,6 +306,34 @@ public final class FleetMcp {
public static LeadConfigDirSource none() { return new LeadConfigDirSource(_ -> null, _ -> null); }
}
/**
* Pane-discovery facts for {@code fleet_list}'s {@code panes} row — every herdr tab's display
* label, the one deliverability gate the status-gated injector itself reads, whether a
* terminal is bound to a configured architect slot, and whether an observer caller may
* {@code SEND} to it.
*
* @param tabLabels tab id → its display label, read lazily (only once the row is actually
* assembled) since it costs a herdr {@code workspace.list}/{@code tab.list}
* scan; a tab herdr reports with no label maps to a {@code null} value
* @param workspaceLabels workspace id → its display label (the herdr "space" name), read lazily
* the same way as {@code tabLabels}; a workspace herdr reports with no label,
* or one the lookup cannot find, maps to a {@code null} value
* @param deliverable the same gate {@link dev.ltms.fleet.Fleetd#deliverableTo} builds for the
* injector, keyed by terminal id — never a second, separately-derived check
* @param architectSlot the same classifier {@link CallerResolver#boundToArchitectSlot} resolves
* a caller against — never a second, separately-derived check
* @param sendableToObserver the same predicate {@link CallerResolver#sendableObserverTarget}
* builds for the {@code SEND} gate — never a second, separately-derived
* check
*/
public record PaneSource(Supplier<Map<String, String>> tabLabels, Supplier<Map<String, String>> workspaceLabels,
Predicate<String> deliverable, Predicate<String> architectSlot, Predicate<String> sendableToObserver) {
/** Inert source — no labels, no deliverable targets, no architect slots, nothing sendable. */
public static PaneSource none() {
return new PaneSource(Map::of, Map::of, _ -> false, _ -> false, _ -> false);
}
}
/**
* fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this
* daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the
@@ -349,11 +381,9 @@ public final class FleetMcp {
* overload instead of failing to compile, and every existing test — none of which exercises
* {@code Fleetd.main} itself — stays green while the live daemon quietly answers
* {@code NOT_CONFIGURED} to {@code fleet_handover} forever. Collapsing every overload into one
* required-everything constructor turns that mistake into a compile error instead: this
* project's own antidote for a defaulted parameter surviving as an untested decision (see
* {@code FleetdCompletionResolverWiringTest} / {@code FleetdLeadRolloverWiringTest}'s own
* javadoc for the same lesson applied to a different seam). A caller that genuinely wants a
* feature off must now say so explicitly at the call site — {@code null},
* required-everything constructor turns that mistake into a compile error instead, the same
* defaulted-parameter guard this codebase applies to other required wiring. A caller that
* genuinely wants a feature off must now say so explicitly at the call site — {@code null},
* {@link OutageSource#none()}, {@link LeadSeatSource#none()}, {@code List.of()} are all still
* perfectly fine values, just never an implicit default reached by omission.
*
@@ -469,7 +499,10 @@ public final class FleetMcp {
recordPrimarySingleton(primaryRegistry, callerTerminal, caller);
Map<String, Object> a = req.arguments();
String target = str(a, "sessionId");
String content = str(a, "content");
// An observer's SEND reaches a pane that cannot otherwise distinguish this
// from a human paste (see attributeIfObserver); every other caller's content
// passes through unchanged.
String content = attributeIfObserver(caller, str(a, "content"));
String turnId = str(a, "turnId");
String coordId = str(a, "coordId");
if (coordId != null && !coordId.isBlank()) {
@@ -501,8 +534,9 @@ public final class FleetMcp {
? sendAsync(messages, target, content, 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.
// fleet_reply's identity is the CONNECTION, never an argument. The authz check is
// terminal ownership, not a role test: the caller may reply only for its own pane,
// which is why no role appears in the check at all.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> replyHandler =
(exchange, req) -> {
String self = callerTerminal(exchange);
@@ -566,13 +600,20 @@ public final class FleetMcp {
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
if (denied != null) return denied;
// A Supplier: the label lookup costs a herdr scan, and must stay behind
// panesVisible so it only runs for a caller that receives the row at all.
PaneSource panes = new PaneSource(() -> identity.panes().tabLabelsByTabId(),
() -> identity.panes().workspaceLabelsByWorkspaceId(),
Fleetd.deliverableTo(presence, callers::leads, callers::collaborators),
callers::boundToArchitectSlot, callers.sendableObserverTarget());
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
callerTerminal(exchange),
callers.collaborators(), collaboratorsVisibleTo(principal(exchange)),
new CoordinationSource(leadChannel, peers),
coordinatorVisibleTo(principal(exchange)),
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)));
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)),
panes, panesVisibleTo(principal(exchange)), principal(exchange).isObserver());
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -600,7 +641,8 @@ public final class FleetMcp {
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_handover", req.arguments()), null);
if (denied != null) return denied;
return handover(leadRollover, callerTerminal(exchange), req.arguments());
Principal caller = principal(exchange);
return handover(leadRollover, messages, caller.terminal(), caller.ownerKey(), req.arguments());
};
McpSchema.Tool fleetSend = sendTool();
@@ -708,7 +750,8 @@ public final class FleetMcp {
if (!authorizationEnforced) {
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
}
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator())) {
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator(),
callers.sendableObserverTarget())) {
if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
}
@@ -796,6 +839,17 @@ public final class FleetMcp {
return caller.isPrimary() || caller.isArchitect();
}
/**
* Who may see {@code fleet_list}'s {@code panes} array — every role that may
* {@link Authz.Action#SEND} to some other pane. A plain worker holds {@code READ} but never
* {@code SEND}, so it still does not see this array. An observer does hold {@code SEND}, to
* another observer pane only, so it sees the array too — but {@code listFleet} filters its rows
* to {@link CallerResolver#sendableObserverTarget} and reduces each one; see {@code paneRows}.
*/
static boolean panesVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator() || caller.isObserver();
}
/**
* The terminal of the caller on this call's connection, or {@code null} when that caller carries
* no terminal, which is only the unnamed primary. A named lead, an architect and a worker each
@@ -914,6 +968,17 @@ public final class FleetMcp {
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
/**
* The text an observer's {@code SEND} actually delivers: prefixed with the sender's own
* connection-resolved terminal, which the receiving pane cannot otherwise tell apart from a
* human paste. Every other caller's content passes through unchanged. Shared with {@code
* FleetApp}'s REST entry path so both surfaces attribute identically.
*/
public static String attributeIfObserver(Principal caller, String content) {
return caller != null && caller.isObserver()
? "[fleet_send from observer " + caller.terminal() + "]\n" + content : content;
}
/**
* {@code fleet_send}: delegate {@code content} to a worker session and block for its reply.
* The configured profiles are required so a profile name can never bypass target validation.
@@ -1358,9 +1423,10 @@ 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 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.
* {@code turnId} and its ticket id are shown only to the caller whose owner key matches the
* delegation's creator — the unnamed primary matches only a delegation another unnamed
* primary created; 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 callerOwner) {
if (isBlank(sessionId)) {
@@ -1481,14 +1547,15 @@ public final class FleetMcp {
* #474 charter tool-surface gate), so every action here must degrade to a clean, structured
* refusal naming {@code NOT_CONFIGURED} rather than ever throwing.
*/
static McpSchema.CallToolResult handover(LeadRollover leadRollover, String callerTerminal,
static McpSchema.CallToolResult handover(LeadRollover leadRollover, MessageService messages,
String callerTerminal, String callerOwner,
Map<String, Object> args) {
String action = str(args, "action");
if (isBlank(action)) {
return error("action is required: \"open\", \"confirm\", \"cancel\" or \"status\"");
}
return switch (action) {
case "open" -> handoverOpen(leadRollover, callerTerminal, str(args, "reason"));
case "open" -> handoverOpen(leadRollover, messages, callerTerminal, callerOwner, str(args, "reason"));
case "confirm" -> handoverConfirm(leadRollover, callerTerminal, str(args, "token"),
truthy(args, "operatorConfirmed"));
case "cancel" -> handoverCancel(leadRollover, str(args, "token"));
@@ -1504,7 +1571,8 @@ public final class FleetMcp {
* IllegalStateException} for that), both degrade to the same clean {@code NOT_CONFIGURED}
* refusal — never an escaping exception.
*/
private static McpSchema.CallToolResult handoverOpen(LeadRollover leadRollover, String callerTerminal,
private static McpSchema.CallToolResult handoverOpen(LeadRollover leadRollover, MessageService messages,
String callerTerminal, String callerOwner,
String reason) {
if (leadRollover == null) {
return notConfigured();
@@ -1522,6 +1590,9 @@ public final class FleetMcp {
m.put("token", p.token());
m.put("handoverPath", p.handoverPath());
m.put("requestedAtMillis", p.requestedAtMillis());
MessageService.Outstanding outstanding = messages.outstanding(callerOwner);
m.put("outstandingTickets", outstanding.tickets());
m.put("openAsks", outstanding.asks());
return text(json(m));
} catch (IllegalStateException e) {
// leadRollover: was removed from config by a hot reload since this FleetMcp was
@@ -1961,6 +2032,25 @@ public final class FleetMcp {
Map.of(), false, coordination, callerIsPrimary, leadsVisible, membersVisible);
}
/**
* As below, with no pane discovery — {@code panes} is {@link PaneSource#none()} and
* {@code panesVisible} is {@code false}.
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
LoopHealthSource loopHealth,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
Map<String, String> collaborators, boolean collaboratorsVisible,
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, contextGauge, leadConfigDirs, leads, selfTerm, collaborators, collaboratorsVisible,
coordination, callerIsPrimary, leadsVisible, membersVisible, PaneSource.none(), false);
}
/**
* The canonical implementation. {@code contextGauge} is the "lead context gauge" (see
* {@link LeadContextGauge}) — every wrapper overload above passes a freshly constructed one,
@@ -1979,6 +2069,11 @@ public final class FleetMcp {
* {@link #leadsVisibleTo})
* @param membersVisible whether this caller may see the {@code members} array (see
* {@link #membersVisibleTo})
* @param panes pane-discovery facts — labels and the deliverable gate for the
* {@code panes} row; {@link PaneSource#none()} for a caller that
* does not want the row
* @param panesVisible whether this caller may see the {@code panes} array (see
* {@link #panesVisibleTo})
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
@@ -1989,11 +2084,38 @@ public final class FleetMcp {
Map<String, String> leads, String selfTerm,
Map<String, String> collaborators, boolean collaboratorsVisible,
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
boolean leadsVisible, boolean membersVisible,
PaneSource panes, boolean panesVisible) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, contextGauge, leadConfigDirs, leads, selfTerm, collaborators, collaboratorsVisible,
coordination, callerIsPrimary, leadsVisible, membersVisible, panes, panesVisible, false);
}
/**
* As above, plus fleetd #758: an observer sees the {@code panes} array too, but filtered to
* {@link CallerResolver#sendableObserverTarget} and each row reduced to the five fields an
* observer may learn — see {@code paneRows}/{@code paneRow}.
*
* @param callerIsObserver whether the {@code fleet_list} caller is an observer; every wrapper
* overload above passes {@code false}, so a test that wants the
* filtered, reduced view must call this overload with an explicit
* {@code true}
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
LoopHealthSource loopHealth,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
Map<String, String> collaborators, boolean collaboratorsVisible,
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible,
PaneSource panes, boolean panesVisible, boolean callerIsObserver) {
try {
// Neither row's assembly (leadView/memberCapacityView probing herdr for live status)
// runs unless at least one of them needs the live-agent lookup backing it.
Map<String, Agent> live = (leadsVisible || membersVisible)
Map<String, Agent> live = (leadsVisible || membersVisible || panesVisible)
? workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
@@ -2037,6 +2159,10 @@ public final class FleetMcp {
.map(e -> collaboratorRow(e.getKey(), e.getValue()))
.toList());
}
// gate BEFORE assembling the row, so the key is absent rather than present-and-empty.
if (panesVisible) {
result.put("panes", paneRows(live, roster, leads, collaborators, panes, callerIsObserver));
}
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
// key is absent rather than present-and-empty.
@@ -2362,6 +2488,116 @@ public final class FleetMcp {
return m;
}
/**
* {@code panes.tabLabels()}'s herdr scan, or an empty map on a {@code HerdrException} — a
* missing label must not cost the {@code leads}/{@code members}/{@code capacity}/
* {@code coordinator} rows that share {@code listFleet}'s own {@code catch}.
*/
private static Map<String, String> tabLabelsOrEmpty(PaneSource panes) {
try {
return panes.tabLabels().get();
} catch (HerdrException e) {
return Map.of();
}
}
/** As {@link #tabLabelsOrEmpty}, for {@code panes.workspaceLabels()}. */
private static Map<String, String> workspaceLabelsOrEmpty(PaneSource panes) {
try {
return panes.workspaceLabels().get();
} catch (HerdrException e) {
return Map.of();
}
}
/**
* One row per herdr-tracked agent pane, sorted by terminal id for a stable order. {@code live}
* is the same terminal-keyed {@link Agent} map {@code leadView}/{@code memberCapacityView}
* already read, so a pane neither configured as a lead nor spawned as a member — a hand-opened
* tab — still gets a row here.
*
* <p>fleetd #758: for an observer caller ({@code observerView}), the rows are filtered to
* {@link PaneSource#sendableToObserver} before being built, and each row is reduced — see
* {@code paneRow}.
*/
private static List<Map<String, Object>> paneRows(Map<String, Agent> live, List<MemberSession> roster,
Map<String, String> leads, Map<String, String> collaborators, PaneSource panes,
boolean observerView) {
final Map<String, String> tabLabels = tabLabelsOrEmpty(panes);
final Map<String, String> workspaceLabels = workspaceLabelsOrEmpty(panes);
Map<String, MemberSession> byTerminal = roster.stream()
.filter(s -> s.terminalId() != null)
.collect(Collectors.toMap(MemberSession::terminalId, Function.identity(), (_, b) -> b));
return live.values().stream()
.filter(a -> !observerView || panes.sendableToObserver().test(a.terminalId()))
.sorted(Comparator.comparing(Agent::terminalId))
.map(a -> paneRow(a, byTerminal.get(a.terminalId()), leads, collaborators, tabLabels,
workspaceLabels, panes, observerView))
.toList();
}
/**
* @param session the roster entry for this pane's terminal, or {@code null} for a pane the
* daemon never spawned as a member (a hand-opened tab, or a configured lead)
* @param tabLabels tab id → its herdr display label; a tab absent here, or carrying a
* {@code null} label itself, projects as a {@code null} "label"
* @param workspaceLabels workspace id → its herdr display label (the space name); a workspace
* absent here, or carrying a {@code null} label itself, projects as a
* {@code null} "workspaceLabel"
* @param observerView fleetd #758: an observer's row carries only {@code sessionId}, {@code
* label}, {@code status}, {@code role}, {@code deliverable} — never {@code
* paneId} (the {@code fleet_stop} handle), {@code workspaceId},
* {@code workspaceLabel}, {@code tabId}, {@code agentType}, or {@code cwd}
* (a member's worktree path is the lead's business)
*/
private static Map<String, Object> paneRow(Agent a, MemberSession session, Map<String, String> leads,
Map<String, String> collaborators, Map<String, String> tabLabels,
Map<String, String> workspaceLabels, PaneSource panes, boolean observerView) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("sessionId", a.terminalId());
if (!observerView) {
m.put("paneId", a.paneId());
m.put("workspaceId", a.workspaceId());
m.put("workspaceLabel", a.workspaceId() == null ? null : workspaceLabels.get(a.workspaceId()));
m.put("tabId", a.tabId());
}
m.put("label", a.tabId() == null ? null : tabLabels.get(a.tabId()));
if (!observerView) {
m.put("agentType", a.agentType());
}
m.put("status", a.status() == null ? "unknown" : a.status().name().toLowerCase());
m.put("role", paneRole(a.terminalId(), session, leads, collaborators, panes));
m.put("deliverable", panes.deliverable().test(a.terminalId()));
if (!observerView && session != null && session.cwd() != null) {
m.put("cwd", session.cwd());
}
return m;
}
/**
* The role this pane resolves as: a spawned member's own {@link MemberRole}, else "lead" for a
* configured but currently-unoccupied lead pane, else "architect" for a pane bound to a
* configured architect slot with no live member session, else "collaborator" for a configured
* but currently-unoccupied collaborator tab, else "observer" for a pane this daemon neither
* spawned nor configured.
*/
private static String paneRole(String terminal, MemberSession session, Map<String, String> leads,
Map<String, String> collaborators, PaneSource panes) {
if (session != null) {
return session.role().wireName();
}
if (leads.containsKey(terminal)) {
return "lead";
}
if (panes.architectSlot().test(terminal)) {
return "architect";
}
if (collaborators.containsKey(terminal)) {
return "collaborator";
}
return "observer";
}
/**
* Reads the lead context gauge for one lead. {@code configDir} is this lead's configured
* {@code CLAUDE_CONFIG_DIR} override (see {@link LeadConfigDirSource}), derived from
@@ -2567,8 +2803,8 @@ public final class FleetMcp {
+ "orchestrators, each with its sessionId (the address to fleet_send to), "
+ "name, live status, and 'self': true on your own row; this is how you "
+ "discover a peer lead without being told its address. 'members' are the "
+ "sessions delegated to — each with sessionId, paneId, role (architect/dev/"
+ "reviewer), profile (the backend it runs on), state, optional "
+ "sessions delegated to — each with sessionId, paneId, role (" + MemberRole.wireNames()
+ "), profile (the backend it runs on), state, optional "
+ "worktree/branch/owner/agentSessionId, and live herdr status. agentSessionId, "
+ "when present, is the id to pass as fleet_spawn's resumeSessionId to relaunch "
+ "onto that same conversation. It is ABSENT — not a guess — for a member fleetd "
@@ -2652,7 +2888,9 @@ public final class FleetMcp {
+ "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 "
+ "confirming — the response also lists outstandingTickets and openAsks, "
+ "your own async delegations and fleet_ask turns, so their ids can go into "
+ "the handover file too), '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 "
@@ -445,27 +445,6 @@ public final class CompositePeerLauncher implements PeerLauncher {
+ " distinct candidate(s): " + String.join(", ", unreachable));
}
/**
* Refuse an explicit-profile spawn when the profile is at its {@code maxLoad} cap.
*
* <p>maxLoad is a documented, unconditional capacity limit (see {@code FleetConfig.Profile#maxLoad}),
* and the charter makes explicit-profile spawns the normal path — so enforcing it only in placement
* ({@code PlacementPolicyUtil}, package-private, hence not linked) would leave the cap dead config
* on every call that names a profile. Same rule as placement: {@code live >= cap} is at capacity.
*
* <p>Deliberately no fallback to another profile: the caller named {@code profile} for a cost/model
* reason, and silently re-routing a paid-tier (subscription) request elsewhere is worse than
* refusing it. A caller that wants placement should omit the profile and let the policy pick.
*
* <p>Known TOCTOU limitation — documented, not fixed. {@link #liveCount} is read outside any lock and
* {@code SessionManager} registers a session only after {@code launcher.spawn} returns, so two
* genuinely concurrent spawns can both pass this check. The race already exists on the placement
* path. Closing it needs slot reservation in the registry; serializing spawn here would block on
* the readiness gate and is a far worse trade.
*
* @param profile the profile the caller explicitly named
* @throws PlacementException when the profile is at capacity
*/
/**
* Refuse an explicit-profile spawn whose credential is quarantined (CB-578 stage B): a prior
* {@code BACKEND_EXHAUSTED} classification on this profile, or on another profile sharing its
@@ -527,6 +506,27 @@ public final class CompositePeerLauncher implements PeerLauncher {
.collect(Collectors.toSet());
}
/**
* Refuse an explicit-profile spawn when the profile is at its {@code maxLoad} cap.
*
* <p>maxLoad is a documented, unconditional capacity limit (see {@code FleetConfig.Profile#maxLoad}),
* and the charter makes explicit-profile spawns the normal path — so enforcing it only in placement
* ({@code PlacementPolicyUtil}, package-private, hence not linked) would leave the cap dead config
* on every call that names a profile. Same rule as placement: {@code live >= cap} is at capacity.
*
* <p>No fallback to another profile: the caller named {@code profile} for a cost/model
* reason, and silently re-routing a paid-tier (subscription) request elsewhere is worse than
* refusing it. A caller that wants placement should omit the profile and let the policy pick.
*
* <p>Known TOCTOU limitation — documented, not fixed. {@link #liveCount} is read outside any lock and
* {@code SessionManager} registers a session only after {@code launcher.spawn} returns, so two
* genuinely concurrent spawns can both pass this check. The race already exists on the placement
* path. Closing it needs slot reservation in the registry; serializing spawn here would block on
* the readiness gate and is a far worse trade.
*
* @param profile the profile the caller explicitly named
* @throws PlacementException when the profile is at capacity
*/
private void enforceMaxLoad(String profile) {
// Absent config, or a config whose maxLoad normalized to null (ABSENT ⇒ unlimited at load),
// means no cap — never cap what wasn't configured. Note "non-positive ⇒ unlimited" was true
@@ -474,7 +474,6 @@ public final class EnvAllowListScrub {
}
}
/** Best-effort recursive delete; failures are swallowed — JVM-exit cleanup is the backstop. */
/**
* Remove generated directories left behind by an earlier daemon process.
*
@@ -511,6 +510,7 @@ public final class EnvAllowListScrub {
}
}
/** Best-effort recursive delete; failures are swallowed — JVM-exit cleanup is the backstop. */
static void deleteRecursively(Path dir) {
if (dir == null || !Files.exists(dir)) {
return;
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.PromptBox;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -23,7 +24,8 @@ import java.util.function.Supplier;
* <p><strong>Status-gated, exactly like {@link ReplyPushLoop}.</strong> A pane may only be injected
* into at a turn boundary ({@link AgentStatus#injectable()} — idle, blocked or done); pasting into
* a live turn corrupts it. So a tick that finds the lead busy simply does nothing and comes back
* later.
* later. The same holds for a lead whose prompt box holds unsubmitted text ({@link PromptBox}) —
* delivering there would submit the operator's half-typed line along with the message.
*
* <p><strong>Ack only after delivery.</strong> A message is acked — removed from the broker — only
* once {@link AgentControl#send} has actually put it in the pane. Anything not delivered (no lead
@@ -52,6 +54,7 @@ public final class LeadCoordLoop {
private final LeadChannel channel;
private final AgentControl agents;
private final PromptBox promptBox;
private final Supplier<Map<String, String>> leads;
private final ScheduledExecutorService scheduler;
private final long intervalMs;
@@ -73,6 +76,7 @@ public final class LeadCoordLoop {
ScheduledExecutorService scheduler, long intervalMs) {
this.channel = channel;
this.agents = agents;
this.promptBox = new PromptBox(agents);
this.leads = leads;
this.scheduler = scheduler;
this.intervalMs = intervalMs;
@@ -149,6 +153,11 @@ public final class LeadCoordLoop {
lead, status, held.size());
return;
}
if (!promptBox.clearToSubmit(lead)) {
log.debug("lead coordination: lead {} has unsubmitted text in its prompt box, holding {} message(s)",
lead, held.size());
return;
}
try {
agents.send(lead, DELIVERY_FORMAT.formatted(msg.from(), msg.content()));
} catch (RuntimeException e) {
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.PromptBox;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.metrics.FleetMetrics;
@@ -33,7 +34,7 @@ import java.util.function.Supplier;
* no such block it is never constructed, so upgrading the daemon cannot silently acquire a behaviour
* that spends the operator's model subscription on its own initiative (constraint 1).
*
* <p>Four invariants keep it from becoming a runaway subscription burner:
* <p>Five invariants keep it from becoming a runaway subscription burner:
* <ol>
* <li><b>Status-gated</b> — a {@code WORKING} lead is making progress and is never touched; only an
* injectable (idle/done/blocked) lead is even considered (constraint 2).</li>
@@ -45,6 +46,8 @@ import java.util.function.Supplier;
* <li><b>Never races {@link ReplyPushLoop}</b> — while that loop is actively nudging any target this
* loop stands down, so two competing injections never start two turns in the same pane
* (constraint 6).</li>
* <li><b>Never submits the operator's draft</b> — a nudge is held while the lead's prompt box holds
* unsubmitted text ({@link PromptBox}), because the delivery pastes and submits in one call.</li>
* </ol>
*
* <p><b>fleetd #609 — context-high notice.</b> Optionally ({@code contextHighNudge}, opt-in like the
@@ -65,6 +68,7 @@ public final class LeadHeartbeatLoop {
private final PrimaryRegistry primaryRegistry;
private final AgentControl agents;
private final PromptBox promptBox;
private final ReplyInbox inbox;
private final Supplier<List<MemberSession>> roster;
private final ReplyPushLoop pushLoop;
@@ -135,6 +139,7 @@ public final class LeadHeartbeatLoop {
boolean requireOperatorConfirm) {
this.primaryRegistry = primaryRegistry;
this.agents = agents;
this.promptBox = new PromptBox(agents);
this.inbox = inbox;
this.roster = roster;
this.pushLoop = pushLoop;
@@ -187,7 +192,9 @@ public final class LeadHeartbeatLoop {
/** Idle past the quiet period with nothing pending and the cap exhausted — stop until new state appears. */
QUIET_DONE,
/** {@link ReplyPushLoop} is actively nudging — stand aside rather than start a competing turn. */
STAND_DOWN
STAND_DOWN,
/** The lead's prompt box holds unsubmitted text — hold the nudge rather than submit that text. */
DRAFT_HELD
}
/**
@@ -330,6 +337,7 @@ public final class LeadHeartbeatLoop {
idleSinceNanos == NOT_IDLE ? null : idleSinceNanos,
quietCount, status, pushLoop.isActive(), leadKnown, fleet,
reading.state(), contextNotified);
d = holdIfOperatorIsTyping(d);
applyDecision(d);
switch (d.action()) {
case INJECT -> injectNudge(d, fleet, reading);
@@ -337,11 +345,33 @@ public final class LeadHeartbeatLoop {
countNudge("exhausted");
contextNotified = d.contextNotified();
}
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> contextNotified = d.contextNotified();
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN, DRAFT_HELD -> contextNotified = d.contextNotified();
}
scheduleNext();
}
/**
* Turn a decision to inject into {@link Action#DRAFT_HELD} when the lead's prompt box holds text
* the operator has not submitted. The pane read happens only for a decision that would otherwise
* send, so a busy or debouncing lead costs no extra herdr call.
*
* <p>The held decision carries this tick's idle window but the <em>pre-tick</em> quiet count and
* context latch: nothing reached the pane, so neither the quiet budget nor the one context notice
* per stretch may be spent on it.
*/
private Decision holdIfOperatorIsTyping(Decision d) {
if (d.action() != Action.INJECT) {
return d;
}
var lead = primaryRegistry.currentPrimaryTerminal();
if (lead.isEmpty() || promptBox.clearToSubmit(lead.get())) {
return d;
}
log.debug("idle-heartbeat: lead {} has unsubmitted text in its prompt box, holding the nudge",
lead.get());
return new Decision(Action.DRAFT_HELD, d.idleSinceNanos(), quietCount, contextNotified);
}
/**
* Persist the idle/quiet state a decision returned, so the next tick starts from it.
*
@@ -12,6 +12,7 @@ import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
@@ -1391,21 +1392,23 @@ public final class MessageService {
}
/**
* 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.
* As {@link #poll(String, String)}, but bypasses the ownership check entirely via
* {@link #INTERNAL_NO_OWNER_CHECK}. No production code calls this overload — it exists for
* tests that only need the ticket's state and have no caller identity to pass.
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
return poll(ticket, INTERNAL_NO_OWNER_CHECK);
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
* 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.
* carries no reply text. The unnamed primary's owner key is {@code null}, matched the same way
* as any other key — it reads a ticket another unnamed primary created, and is refused on a
* ticket a named caller created. 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 callerOwner) {
Task task = tasks.get(ticket);
@@ -1452,13 +1455,27 @@ public final class MessageService {
}
/**
* 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
* Marker passed as {@code callerOwner} to bypass the ownership check entirely. No
* {@link Principal#ownerKey()} ever produces this value — every real key is either
* {@code null} (the unnamed primary) or prefixed with its role, such as {@code "worker:"} or
* {@code "leader:"}. {@link #poll(String)} passes it; {@link #pendingAsk} has no matching
* no-check overload, so this stays package-private for the test that drives the bypass
* directly.
*/
static final String INTERNAL_NO_OWNER_CHECK = "internal:no-owner-check";
/**
* Whether {@code callerOwner} may read {@code task}'s state. {@code callerOwner} is matched
* against the task's recorded owner key by equality, including a {@code null} match — the
* unnamed primary's owner key is {@code null}, so it owns a ticket another unnamed primary
* created and nothing else, the same rule every other role follows. The only caller that
* reads any ticket is {@link #INTERNAL_NO_OWNER_CHECK}. 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 callerOwner) {
return callerOwner == null || callerOwner.equals(task.creatorOwner);
return INTERNAL_NO_OWNER_CHECK.equals(callerOwner)
|| Objects.equals(callerOwner, task.creatorOwner);
}
/**
@@ -1797,6 +1814,67 @@ public final class MessageService {
return null;
}
/**
* One ticket {@code callerOwner} created, still present in {@link #tasks}, surfaced by
* {@link #outstanding} so a lead can carry its id into a handover file. {@link #phase} is the
* same value {@link #poll} would report right now, terminal phases included: a {@code DONE} or
* {@code FAILED} ticket stays in {@link #tasks} — and so stays reported here — until
* {@link #pruneTerminalTickets} evicts it.
*/
public record OutstandingTicket(String ticket, Phase phase, String target) {
}
/**
* One worker session paused in {@code fleet_ask}, with the {@code turnId} that answers it,
* surfaced by {@link #outstanding} alongside {@link OutstandingTicket}.
*/
public record OutstandingAsk(String ticket, String turnId, String workerSession) {
}
/** The outstanding tickets and open asks a single call to {@link #outstanding} reports. */
public record Outstanding(List<OutstandingTicket> tickets, List<OutstandingAsk> asks) {
}
/**
* Every ticket {@code callerOwner} created that is still in {@link #tasks} — including a
* finished one nobody has polled yet, since {@link #pruneTerminalTickets} discards its reply
* on a timer and a lead that does not carry its id forward can no longer read it after losing
* its session's context — plus the subset of those whose worker is paused in
* {@code fleet_ask}. Filtered by the same ownership rule as {@link #poll}:
* {@link #ownsTicket(Task, String)}.
*/
public Outstanding outstanding(String callerOwner) {
List<OutstandingTicket> tickets = new ArrayList<>();
List<OutstandingAsk> asks = new ArrayList<>();
for (Task task : tasks.values()) {
if (!ownsTicket(task, callerOwner)) {
continue;
}
Reply question = task.question;
Phase phase;
if (task.future.isDone()) {
phase = terminalPhase(task.future);
} else if (question != null) {
phase = Phase.ASKING;
asks.add(new OutstandingAsk(task.ticket, question.turnId(), task.target));
} else {
phase = Phase.PENDING;
}
tickets.add(new OutstandingTicket(task.ticket, phase, task.target));
}
return new Outstanding(tickets, asks);
}
/** As {@link #poll}'s own terminal-result handling, reduced to just the {@link Phase}. */
private static Phase terminalPhase(CompletableFuture<Reply> future) {
try {
Reply r = future.getNow(null);
return r != null && r.completed() ? Phase.DONE : Phase.FAILED;
} catch (CompletionException | java.util.concurrent.CancellationException e) {
return Phase.FAILED;
}
}
/** Release the async executor. */
public void close() {
asyncExecutor.shutdown();
@@ -3,6 +3,7 @@ package dev.ltms.fleet.msg;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.PromptBox;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
@@ -50,6 +51,11 @@ import java.util.stream.Collectors;
* exhausting its cap does not stop nudges about the others (post-CB-590 regression fix; see
* {@link #decide}) — whichever the durable inbox / pending set doesn't already answer via
* {@code STOP}.
*
* <p>A lead that is injectable is nudged only when its prompt box is also empty
* ({@link PromptBox}): the delivery pastes and submits in one call, so a nudge into a box holding
* the operator's half-typed line would submit that line too. A nudge held for that reason waits for
* the next tick like any other, and the pending work is re-read then.
*/
public final class ReplyPushLoop {
@@ -81,6 +87,7 @@ public final class ReplyPushLoop {
private final PrimaryRegistry primaryRegistry;
private final AgentControl agents;
private final PromptBox promptBox;
private final ReplyInbox inbox;
private final ScheduledExecutorService scheduler;
private final int maxReminders;
@@ -125,6 +132,7 @@ public final class ReplyPushLoop {
int maxReminders, long backoffMs, Metrics metrics) {
this.primaryRegistry = primaryRegistry;
this.agents = agents;
this.promptBox = new PromptBox(agents);
this.inbox = inbox;
this.scheduler = scheduler;
this.maxReminders = maxReminders;
@@ -398,11 +406,15 @@ public final class ReplyPushLoop {
log.debug("push: status check failed for lead {}, will retry", lead, e);
return Action.WAIT_BUSY;
}
if (status.injectable()) {
return Action.INJECT;
if (!status.injectable()) {
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
return Action.WAIT_BUSY;
}
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
return Action.WAIT_BUSY;
if (!promptBox.clearToSubmit(lead)) {
log.debug("push: lead {} has unsubmitted text in its prompt box, waiting", lead);
return Action.WAIT_BUSY;
}
return Action.INJECT;
}
/**
@@ -1,6 +1,8 @@
package dev.ltms.fleet.peer;
import java.util.Locale;
import java.util.stream.Collectors;
import java.util.stream.Stream;
/**
* What a member is <em>for</em> — the contract it runs under.
@@ -100,6 +102,11 @@ public enum MemberRole {
return null;
}
/** The wire name of every role, joined with {@code ", "} in declaration order. */
public static String wireNames() {
return Stream.of(values()).map(MemberRole::wireName).collect(Collectors.joining(", "));
}
/**
* Parse a config/wire spelling, case-insensitively.
*
@@ -118,14 +125,7 @@ public enum MemberRole {
}
}
}
StringBuilder valid = new StringBuilder();
for (MemberRole r : values()) {
if (!valid.isEmpty()) {
valid.append(", ");
}
valid.append(r.wireName());
}
throw new IllegalArgumentException(
"unknown member role '" + s + "'; valid roles are: " + valid);
"unknown member role '" + s + "'; valid roles are: " + wireNames());
}
}
@@ -13,6 +13,7 @@ import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.member.MemberCredentialPolicyView;
import dev.ltms.fleet.peer.PeerUnreachableException;
@@ -281,6 +282,18 @@ public final class FleetApp {
return Authz.permits(caller, action, target, knownLeadOrCollaborator);
}
/**
* As {@link #permitsFor(Principal, Authz.Action, String, Predicate)}, also threading the
* classifier an observer's {@code SEND} is checked against; pass {@link #auth}'s own
* {@code sendableObserverTarget()} to exercise the real production gate, as {@link #allow}
* does.
*/
static boolean permitsFor(Principal caller, Authz.Action action, String target,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> knownObserverTarget) {
return Authz.permits(caller, action, target, knownLeadOrCollaborator, knownObserverTarget);
}
/**
* Gate a handler on the CB-505 authorization table. Returns {@code true} when the request may
* proceed; otherwise writes the error response and returns {@code false}.
@@ -294,7 +307,7 @@ public final class FleetApp {
return true; // legacy: authorization not enforced
}
Principal caller = ctx.attribute(CALLER);
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator())) {
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.sendableObserverTarget())) {
if (action != Authz.Action.READ && action != Authz.Action.METRICS
&& action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
@@ -436,14 +449,27 @@ public final class FleetApp {
}
}
/**
* Tab id → its herdr display label, or an empty map on a {@code workspace.list}/{@code
* tab.list} failure — a missing label must not cost the agent roster.
*/
private Map<String, String> tabLabelsOrEmpty() {
try {
return new PaneLocator(herdr, memberHerdr).tabLabelsByTabId();
} catch (HerdrException e) {
return Map.of();
}
}
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
private void agents(Context ctx) {
if (!allow(ctx, routeAction("GET /agents"), null)) {
return;
}
final Map<String, String> tabLabels = tabLabelsOrEmpty();
try {
ctx.status(200).json(Map.of("agents",
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
workers.list().stream().map(Agent.class::cast).map(a -> view(a, tabLabels)).toList()));
} catch (HerdrException e) {
// fleetd #297: workers.list() reaches herdr — a transport failure must land in the same
// {error, detail} envelope every other failure path here uses, not escape as a bare
@@ -686,6 +712,10 @@ public final class FleetApp {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "content is required"));
return;
}
// An observer's SEND reaches a pane that cannot otherwise distinguish this from a human
// paste (see FleetMcp#attributeIfObserver, the same rule on the MCP entry path); every
// other caller's content passes through unchanged.
content = FleetMcp.attributeIfObserver(caller, content);
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
@@ -883,8 +913,7 @@ 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 owner key created that delegation, or
// to the unnamed primary.
// Phase.ASKING view, but only to the caller whose owner key created that delegation.
Principal caller = ctx.attribute(CALLER);
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.ownerKey());
if (ask != null) {
@@ -936,13 +965,19 @@ public final class FleetApp {
}
}
/** Stable JSON projection of an agent (null-safe for the start-time shape). */
private static Map<String, Object> view(Agent a) {
/**
* Stable JSON projection of an agent (null-safe for the start-time shape).
*
* @param tabLabels tab id → its herdr display label; a tab absent from this map, or carrying
* a {@code null} label itself, projects as a {@code null} "label"
*/
private static Map<String, Object> view(Agent a, Map<String, String> tabLabels) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("terminalId", a.terminalId());
m.put("paneId", a.paneId());
m.put("workspaceId", a.workspaceId());
m.put("tabId", a.tabId());
m.put("label", tabLabels.get(a.tabId()));
m.put("sessionId", a.sessionId());
m.put("agentType", a.agentType());
m.put("status", a.status().name().toLowerCase());
@@ -442,38 +442,6 @@ public final class GitWorktrees implements Worktrees {
ENVIRONMENT_CREDENTIAL_HELPER);
}
/**
* {@link #configureEnvironmentCredentialHelper} only ever fires for an HTTPS origin — Git never
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member sitting on that origin never
* reaches the helper and the repo-scoped {@code WORKER_GITEA_TOKEN} is simply not used.
*
* <p>An earlier version of this javadoc justified the rewrite by claiming a member <em>cannot</em>
* push once {@code memberCredentials.policy: allow-list} blocks {@code SSH_AUTH_SOCK}, because
* "there is no private key file on this host, only an ssh-agent socket". That premise is false
* (fleetd #184): {@code ssh -G} resolves a readable, passphrase-free {@code IdentityFile} outside
* {@code ~/.ssh}, and a member — same OS user — pushes over SSH with the socket blanked. The
* rewrite is still worth having, but for the reason below rather than that one: it routes the
* member through its own scoped token instead of the operator's ssh identity, which is what makes
* a member's pushes attributable and revocable.
*
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
* <ssh-base>}, set with {@code --worktree} so it lands only in
* {@code <worktree>/.git/worktrees/<name>/config.worktree} (enabled by
* {@code extensions.worktreeConfig}, already turned on above) and never touches the shared
* repo-level config the primary checkout also reads. {@code insteadOf} — not
* {@code pushInsteadOf} — because a member may also need to fetch or rebase, and both should go
* through the member's own token for the same reason.
*
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
* itself — never a hardcoded forge host, which is exactly what #177 removed. An origin that is
* already {@code https://} is left alone; the credential helper already covers it. An origin
* that is neither {@code ssh://} nor {@code https://} — including the scp-like shorthand
* ({@code git@host:path}, no scheme) — is left untouched deliberately: that shorthand's
* {@code host:path} split is defined by the user's ssh_config aliases, not by URI syntax, so
* guessing at it risks rewriting to the wrong place. A repo provisioned from that form keeps
* today's (broken, if the policy blocks the agent) SSH-only behaviour rather than a wrong rewrite.
*/
/**
* Blank the user-info of a remote URL before it reaches a log. A remote URL is not obviously a
* credential channel, which is exactly why one has leaked here three times ({@code git remote -v}
@@ -485,6 +453,30 @@ public final class GitWorktrees implements Worktrees {
return url == null ? null : url.replaceAll("://[^@/]*@", "://<redacted>@");
}
/**
* {@link #configureEnvironmentCredentialHelper} only fires for an HTTPS origin — Git never
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member on that origin never
* reaches the helper, and the repo-scoped token goes unused without a separate rewrite.
*
* <p>This method routes the member through its own scoped token instead of the operator's ssh
* identity, which is what makes a member's pushes attributable and revocable.
*
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
* <ssh-base>}, set with {@code --worktree} so it lands only in
* {@code <worktree>/.git/worktrees/<name>/config.worktree} and never touches the shared
* repo-level config the primary checkout also reads. {@code insteadOf} — not
* {@code pushInsteadOf} — because a member may also need to fetch or rebase through its own
* token.
*
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
* itself, never a hardcoded forge host. An origin already {@code https://} is left alone; the
* credential helper already covers it. An origin that is neither {@code ssh://} nor
* {@code https://} — including the scp-like shorthand ({@code git@host:path}, no scheme) — is
* left untouched: that shorthand's {@code host:path} split is defined by the user's ssh_config
* aliases, not by URI syntax, so guessing at it risks rewriting to the wrong place, and that
* origin keeps SSH-only push behaviour instead.
*/
private void configureHttpsUrlRewriteForSshOrigin(String repoRoot, String worktreePath) {
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
return;
@@ -125,6 +125,7 @@ class FleetdAssemblyTurnRegistrarBehaviouralTest {
primary:
tab: "lead: primary"
profile: sonnet
workspace: "ltms"
profiles:
sonnet:
subscription: true
@@ -91,6 +91,11 @@ class FleetdBackendErrorSinkTest {
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
.put("terminal_id", "term_primary").put("agent_status", "idle"));
}
if ("agent.read".equals(method)) {
// The lead-nudge paths read the input box before pasting into it.
return MAPPER.createObjectNode().set("read",
MAPPER.createObjectNode().put("text", FakeHerdr.IDLE_PROMPT_CARET));
}
if ("agent.prompt".equals(method)) {
prompts.add(params);
sendLatch.countDown();
@@ -171,6 +171,7 @@ class FleetdLeadContextSourceWindowAssemblyTest {
%s:
tab: "%s"
profile: %s
workspace: "ltms"
profiles:
%s:
subscription: true
@@ -172,6 +172,7 @@ class FleetdLeadRolloverAssemblyTest {
opus:
tab: "lead: opus"
cwd: "%s"
workspace: "ltms"
leadRollover:
handoverPath: handover.md
requireOperatorConfirm: false
@@ -203,6 +204,7 @@ class FleetdLeadRolloverAssemblyTest {
tab: "lead: opus"
cwd: "%s"
profile: opus
workspace: "ltms"
profiles:
opus:
subscription: true
@@ -143,6 +143,7 @@ class FleetdLeadSeatAssemblyTest {
opus:
tab: "lead: opus"
profile: sonnet
workspace: "ltms"
profiles:
sonnet:
subscription: true
@@ -1,96 +1,174 @@
package dev.ltms.fleet;
import com.tngtech.archunit.base.DescribedPredicate;
import com.tngtech.archunit.core.domain.Dependency;
import com.tngtech.archunit.core.domain.JavaClass;
import com.tngtech.archunit.core.domain.JavaClass.Predicates;
import com.tngtech.archunit.core.domain.JavaClasses;
import com.tngtech.archunit.core.importer.ClassFileImporter;
import com.tngtech.archunit.core.importer.ImportOption;
import com.tngtech.archunit.library.dependencies.SliceRule;
import com.tngtech.archunit.library.dependencies.SlicesRuleDefinition;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.fail;
/**
* fleetd #131 (CB-627): enforce package boundaries with an ArchUnit test instead of a
* Maven module split.
* Enforces package boundaries between the top-level {@code dev.ltms.fleet.*} packages.
*
* <p>This test fails the build the moment a NEW cycle appears between the top-level
* {@code dev.ltms.fleet.*} packages. Today's cycles are recorded below as explicit,
* narrow exceptions: each one ignores dependencies between exactly the two named
* packages, in both directions, and nothing else. A cycle through any other pair of
* packages -- or a brand new pair -- still fails this test.
* <p>{@link #BASELINE_EDGES} names the exact {@code origin class -> target class}
* dependencies allowed to cross a top-level package boundary. Any dependency between two
* top-level packages that is not in that set fails this test, including a brand new
* dependency between a pair of packages that already has other baselined edges. A baseline
* entry whose dependency no longer exists in the code also fails this test, so the baseline
* always names exactly today's exceptions and nothing more.
*
* <p><b>Main code only.</b> The import excludes test classes
* ({@link ImportOption.Predefined#DO_NOT_INCLUDE_TESTS}). Test code legitimately wires
* across many packages for setup and mocking; that is not part of the shipped
* architecture this rule protects. Verified: importing test classes too pulls in a much
* larger, noisier cycle set -- {@code herdr}, {@code member}, {@code peer}, {@code
* config}, {@code guard} and {@code placement} all show up in cycles that disappear the
* moment test classes are excluded. Scanning off the classpath via {@code
* importPackages(...)} (not a hardcoded {@code target/classes} path) also keeps this
* test correct regardless of the working directory the build is invoked from.
*
* <p><b>No package moves here</b> -- ticket #131 is explicit that removing a cycle is
* its own, later PR. See the comment on each exception below for which ticket step
* removes it.
* ({@link ImportOption.Predefined#DO_NOT_INCLUDE_TESTS}). Scanning off the classpath via
* {@code importPackages(...)} keeps this test correct regardless of the working directory
* the build is invoked from.
*/
class PackageCyclesTest {
/**
* Exact {@code "origin -> target"} class dependencies allowed to cross a top-level
* package boundary. Each entry is one directed edge between two specific classes; a
* two-way relationship between a pair of packages is listed as two separate entries,
* one per direction.
*/
private static final Set<String> BASELINE_EDGES = Set.of(
"dev.ltms.fleet.auth.CallerResolver -> dev.ltms.fleet.mcp.ConnectionIdentity",
"dev.ltms.fleet.auth.CallerResolver -> dev.ltms.fleet.mcp.ConnectionIdentity$Caller",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.AuditLog",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Authz",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Authz$Action",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.CallerResolver",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Principal",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Role",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadChannel",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadChannel$MailboxState",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadMessage",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$AskOutcome",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$AskResult",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Outcome",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Outstanding",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$PendingAsk",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Phase",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Reply",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$ReplyOutcome",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$TaskView",
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$AskOutcome",
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$Outcome",
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$Phase",
"dev.ltms.fleet.mcp.FleetMcp$CoordinationSource -> dev.ltms.fleet.msg.LeadChannel",
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.mcp.PrimaryRegistry",
"dev.ltms.fleet.msg.ReplyPushLoop -> dev.ltms.fleet.mcp.PrimaryRegistry",
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.Rendezvous",
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.Rendezvous$Resolution",
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.inject.CompletionResolver$InFlight -> dev.ltms.fleet.msg.Rendezvous$Resolution",
"dev.ltms.fleet.inject.Injector -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.inject.Injector$Pending -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.inject.TurnListener -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.inject.TurnRegistrar -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector",
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector$Cancellation",
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector$Delivery",
"dev.ltms.fleet.metrics.FleetMetrics -> dev.ltms.fleet.msg.ReplyInbox",
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.metrics.Metrics",
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.metrics.Metrics",
"dev.ltms.fleet.msg.ReplyPushLoop -> dev.ltms.fleet.metrics.Metrics",
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.session.MemberSession",
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.session.MemberSession$State",
"dev.ltms.fleet.session.SessionManager -> dev.ltms.fleet.msg.TurnToken"
);
private static final String ROOT_PACKAGE = "dev.ltms.fleet.";
@Test
void packagesAreFreeOfCycles() {
var classes = new ClassFileImporter()
JavaClasses classes = new ClassFileImporter()
.withImportOption(ImportOption.Predefined.DO_NOT_INCLUDE_TESTS)
.importPackages("dev.ltms.fleet");
checkBaselineMatchesTodaysEdges(classes);
SliceRule rule = SlicesRuleDefinition.slices()
.matching("dev.ltms.fleet.(*)..")
.should().beFreeOfCycles();
// fleetd #131 step 1: move ConnectionIdentity so authz stops depending on the
// MCP layer. Evidence: auth/CallerResolver.java:3 imports mcp.ConnectionIdentity;
// mcp/FleetMcp.java:3-7 imports auth.AuditLog, Authz, CallerResolver, Principal,
// Role.
rule = ignoreCycle(rule, "auth", "mcp");
// fleetd #131 step 2: PrimaryRegistry is used by loops in msg; move it, or put
// an interface between msg and mcp. Evidence: msg/ReplyPushLoop.java:5 and
// msg/LeadHeartbeatLoop.java:5 import mcp.PrimaryRegistry; mcp/FleetMcp.java:15-18
// imports msg.LeadChannel, LeadMessage, MessageService, Rendezvous.
rule = ignoreCycle(rule, "mcp", "msg");
// fleetd #131 -- found while implementing this test, NOT one of the ticket's
// original three; it names its own follow-up step before removal. Evidence:
// inject/CompletionResolver.java:4-5, inject/Injector.java:6 and
// inject/TurnListener.java:3 import msg.Rendezvous / msg.TurnToken;
// msg/MessageService.java:6 imports inject.Injector.
rule = ignoreCycle(rule, "inject", "msg");
// fleetd #131 -- same as above, its own follow-up. Evidence:
// metrics/FleetMetrics.java:3 imports msg.ReplyInbox; msg/MessageService.java:7-8,
// msg/LeadHeartbeatLoop.java:6-7 and msg/ReplyPushLoop.java:6-7 import
// metrics.FleetMetrics / metrics.Metrics.
rule = ignoreCycle(rule, "metrics", "msg");
// fleetd #131 -- same as above, its own follow-up. Evidence:
// session/SessionManager.java:7 imports msg.TurnToken;
// msg/LeadHeartbeatLoop.java:8 imports session.MemberSession.
rule = ignoreCycle(rule, "msg", "session");
for (String edge : BASELINE_EDGES) {
String[] originAndTarget = edge.split(" -> ");
rule = rule.ignoreDependency(originAndTarget[0], originAndTarget[1]);
}
rule.check(classes);
}
/**
* Accepts today's known cycle between two top-level packages, and nothing else.
* Ignoring both directions removes exactly this pair from cycle detection; every
* other dependency -- including any new one added later, between these same two
* packages or any other pair -- is still checked.
* Fails with the exact offending edge when the live code and {@link #BASELINE_EDGES}
* disagree: a dependency crossing a baselined package pair that is not in the baseline,
* or a baseline entry whose dependency no longer exists.
*/
private static SliceRule ignoreCycle(SliceRule rule, String packageA, String packageB) {
return rule
.ignoreDependency(residesIn(packageA), residesIn(packageB))
.ignoreDependency(residesIn(packageB), residesIn(packageA));
private static void checkBaselineMatchesTodaysEdges(JavaClasses classes) {
Set<String> baselinedPackagePairs = new TreeSet<>();
for (String edge : BASELINE_EDGES) {
String[] originAndTarget = edge.split(" -> ");
baselinedPackagePairs.add(unorderedPair(
topLevelPackageOf(originAndTarget[0]), topLevelPackageOf(originAndTarget[1])));
}
Set<String> liveEdgesInBaselinedPairs = new TreeSet<>();
for (JavaClass javaClass : classes) {
for (Dependency dependency : javaClass.getDirectDependenciesFromSelf()) {
JavaClass origin = dependency.getOriginClass();
JavaClass target = dependency.getTargetClass();
String originPackage = topLevelPackageOf(origin.getFullName());
String targetPackage = topLevelPackageOf(target.getFullName());
if (originPackage.isEmpty() || targetPackage.isEmpty() || originPackage.equals(targetPackage)) {
continue;
}
if (baselinedPackagePairs.contains(unorderedPair(originPackage, targetPackage))) {
liveEdgesInBaselinedPairs.add(origin.getFullName() + " -> " + target.getFullName());
}
}
}
List<String> problems = new ArrayList<>();
for (String liveEdge : liveEdgesInBaselinedPairs) {
if (!BASELINE_EDGES.contains(liveEdge)) {
String[] originAndTarget = liveEdge.split(" -> ");
problems.add("new dependency not in the baseline: " + liveEdge
+ " (packages " + topLevelPackageOf(originAndTarget[0])
+ " -> " + topLevelPackageOf(originAndTarget[1]) + ")");
}
}
for (String baselineEdge : BASELINE_EDGES) {
if (!liveEdgesInBaselinedPairs.contains(baselineEdge)) {
String[] originAndTarget = baselineEdge.split(" -> ");
problems.add("stale baseline entry, no such dependency exists: " + baselineEdge
+ " (packages " + topLevelPackageOf(originAndTarget[0])
+ " -> " + topLevelPackageOf(originAndTarget[1]) + ")");
}
}
if (!problems.isEmpty()) {
fail("PackageCyclesTest baseline is out of date:\n " + String.join("\n ", problems));
}
}
private static DescribedPredicate<JavaClass> residesIn(String topLevelPackage) {
return Predicates.resideInAPackage("dev.ltms.fleet." + topLevelPackage + "..");
private static String unorderedPair(String packageA, String packageB) {
return packageA.compareTo(packageB) <= 0 ? packageA + "|" + packageB : packageB + "|" + packageA;
}
private static String topLevelPackageOf(String fullyQualifiedClassName) {
if (!fullyQualifiedClassName.startsWith(ROOT_PACKAGE)) {
return "";
}
String rest = fullyQualifiedClassName.substring(ROOT_PACKAGE.length());
int dot = rest.indexOf('.');
return dot < 0 ? "" : rest.substring(0, dot);
}
}
@@ -288,14 +288,16 @@ class AuthzTest {
}
/**
* 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.
* Every action beyond READ/METRICS/REPLY/ASK/SEND, 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. {@code SEND} is excluded
* here and given its own matrix below, since — unlike every action in this loop — its grant is
* conditional on the target, not fixed.
*/
@Test
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAndAsk() {
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskAndSend() {
for (Authz.Action a : Authz.Action.values()) {
if (a == READ || a == METRICS || a == REPLY || a == ASK) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == SEND) {
continue;
}
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
@@ -312,4 +314,53 @@ class AuthzTest {
assertFalse(OBSERVER.isSpawnedMember());
assertTrue(OBSERVER.isObserver());
}
// ── the observer SEND matrix ────────────────────────────────────────────────────────────────
/**
* {@code SEND} for an observer is the one grant that is conditional rather than fixed, exactly
* like a collaborator's: flipping only the observer-target classifier's answer flips only this
* outcome.
*/
@Test
void anObserverMaySendOnlyWhenTheClassifierAcceptsTheTargetAsAnObserver() {
assertTrue(Authz.permits(OBSERVER, SEND, "term_other_observer",
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, target -> true),
"the classifier accepting the target as an observer must grant SEND");
assertFalse(Authz.permits(OBSERVER, SEND, "term_other_observer",
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, target -> false),
"the classifier refusing the target must deny SEND");
assertFalse(Authz.permits(OBSERVER, SEND, "term_other_observer"),
"the real production classifier recognises no terminal as an observer target yet, "
+ "so SEND is refused by the three- and four-argument convenience forms");
}
/**
* Control for the test above: every other action's result for an observer does not move when
* the observer-target classifier does. Only {@code SEND} is wired to it.
*/
@Test
void theObserverTargetClassifierMovesOnlySendForAnObserver() {
for (Authz.Action a : Authz.Action.values()) {
if (a == SEND) {
continue;
}
assertEquals(
Authz.permits(OBSERVER, a, "term_observer"),
Authz.permits(OBSERVER, a, "term_observer",
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, target -> true),
a + " must not depend on the observer-target classifier at all");
}
}
/**
* An observer's {@code SEND} is gated on a different classifier than a collaborator's: the
* collaborator classifier accepting every target must not itself grant an observer's SEND.
*/
@Test
void anObserversSendDoesNotMoveOnTheCollaboratorClassifier() {
assertFalse(Authz.permits(OBSERVER, SEND, "term_lead", target -> true),
"an observer's SEND must consult the observer-target classifier, never the "
+ "lead-or-collaborator one");
}
}
@@ -912,4 +912,72 @@ class CallerResolverTest {
assertFalse(r.knownLeadOrCollaborator().test("term_a"));
}
// ── fleetd #743: sendableObserverTarget() reads the same maps and functions resolve() does ────
/**
* A terminal this resolver recognises as none of the privileged roles is exactly the one
* {@code resolve} would itself hand back {@link Role#OBSERVER} for.
*/
@Test
void sendableObserverTargetIsTrueForATerminalKnownAsNoOtherRole() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_lead", "opus-5.0"), new MemberRegistry(null),
t -> "term_worker".equals(t) ? MemberRole.DEV : null,
() -> Map.of("term_collab", "ops"));
assertTrue(r.sendableObserverTarget().test("term_other"));
}
@Test
void sendableObserverTargetIsFalseForALeadTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_lead", "opus-5.0"), new MemberRegistry(null), t -> null, Map::of);
assertFalse(r.sendableObserverTarget().test("term_lead"),
"a lead's own terminal must never be a sendable observer target");
}
@Test
void sendableObserverTargetIsFalseForACollaboratorTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, () -> Map.of("term_collab", "ops"));
assertFalse(r.sendableObserverTarget().test("term_collab"),
"a collaborator's own terminal must never be a sendable observer target");
}
@Test
void sendableObserverTargetIsFalseForALiveSpawnedMembersTerminal() {
// Covers both a worker and an architect: spawnedMemberRole.apply(target) is non-null for
// either, and resolve() never falls through to OBSERVER once it is.
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null),
t -> switch (t) {
case "term_worker" -> MemberRole.DEV;
case "term_architect" -> MemberRole.ARCHITECT;
default -> null;
}, Map::of);
assertFalse(r.sendableObserverTarget().test("term_worker"));
assertFalse(r.sendableObserverTarget().test("term_architect"));
}
@Test
void sendableObserverTargetIsFalseForABoundArchitectSlotWithNoLiveMember() {
// The edge case resolve() itself carries: a terminal bound to a configured architect slot
// but with no live spawned-member session yet.
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
boundMembers("architect:lead-designer", MemberRole.ARCHITECT), t -> null, Map::of);
assertFalse(r.sendableObserverTarget().test("term_a"));
}
@Test
void sendableObserverTargetIsFalseForANullTarget() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, Map::of);
assertFalse(r.sendableObserverTarget().test(null));
}
}
@@ -537,12 +537,13 @@ class FleetConfigTest {
}
/**
* CB-579: {@code tab} is the only field a lead's identity depends on now, so it is required
* whether the entry is creatable or recognise-only — without it the entry can never be found.
* A lead's tab label is a fixed constant, not a per-entry field, so an entry with no {@code
* tab:} is the normal case — it is still found by that constant label in its own {@code
* workspace}, not refused as useless.
*/
@Test
void aLeadWithNoTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("useless-lead.yaml");
void aLeadWithNoTabIsAcceptedAndFoundByTheFixedLabel(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-tab-lead.yaml");
Files.writeString(f, """
bind:
port: 8080
@@ -553,9 +554,9 @@ class FleetConfigTest {
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateMembers);
assertTrue(e.getMessage().contains("ghost"), "the message must name the useless entry");
assertTrue(e.getMessage().contains("tab:"), "the message must say what is missing");
assertDoesNotThrow(cfg::validateMembers);
assertEquals(List.of(FleetConfig.Leader.LEAD_TAB_LABEL),
cfg.fleet().leaders().get("ghost").acceptedLabels());
}
/**
@@ -775,22 +776,45 @@ class FleetConfigTest {
}
/**
* fleetd #677: identity is matched on a lead's exact {@code tab} alone, so two leads sharing
* one tab means only one of them is ever found — the guard must catch this independently of
* the member-template checks above.
* A lead's tab label is fixed, so two leads sharing one {@code workspace} would both resolve
* to the one tab named {@code lead} there — the guard must catch this independently of the
* member-template checks above.
*/
@Test
void twoLeadsSharingTheSameExactTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-tab.yaml");
void twoLeadersSharingTheSameWorkspaceRefuseToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-workspace.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "shared tab"
workspace: "shared"
sonnet:
tab: "shared tab"
workspace: "shared"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("opus"), "the message must name one offending lead");
assertTrue(e.getMessage().contains("sonnet"), "the message must name the other offending lead");
assertTrue(e.getMessage().contains("shared"), "the message must name the shared workspace");
}
/** The workspace collision check is case-insensitive, matching how spaces are looked up. */
@Test
void twoLeadersSharingTheSameWorkspaceInDifferentCaseRefuseToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-workspace-case.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
workspace: "Shared"
sonnet:
workspace: "shared"
""");
FleetConfig cfg = FleetConfig.load(f);
@@ -800,52 +824,77 @@ class FleetConfigTest {
assertTrue(e.getMessage().contains("sonnet"), "the message must name the other offending lead");
}
/**
* fleetd #693: the guard matches tabs case-insensitively, because
* {@code LeadTabScanner} keys its tab map on a lowercased label — two tabs differing only in
* case collide there too, and the guard must catch that independently of the exact-match case
* above.
*/
/** Control for the two tests above: distinct workspaces load cleanly, with no {@code tab:} at all. */
@Test
void twoLeadsSharingTheSameTabInDifferentCaseRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-tab-case.yaml");
void twoLeadersWithDistinctWorkspacesAndNoTabAreAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("distinct-workspaces.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "Shared Tab"
workspace: "space-opus"
sonnet:
tab: "shared tab"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("opus"), "the message must name one offending lead");
assertTrue(e.getMessage().contains("sonnet"), "the message must name the other offending lead");
}
/** Control for {@link #twoLeadsSharingTheSameExactTabRefusesToStart}: distinct tabs load cleanly. */
@Test
void twoLeadsWithDistinctExactTabsAreAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("distinct-tabs.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "opus tab"
sonnet:
tab: "sonnet tab"
workspace: "space-sonnet"
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateLeadTabPrefixes);
}
/**
* A member tab-label template that can render exactly as the fixed lead tab label would let a
* member's own tab be read back as a lead — refused outright, with no lead needing to be
* configured at all.
*/
@Test
void aFleetTabLabelTemplateThatCanRenderAsTheFixedLeadTabLabelRefusesToStart(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("template-renders-as-lead.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
tabLabel: "lead"
leaders:
opus:
workspace: "fleet"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("fleet.tabLabel"),
"the message must name the offending template");
assertTrue(e.getMessage().contains(FleetConfig.Leader.LEAD_TAB_LABEL),
"the message must name the fixed lead tab label it collides with");
}
/**
* A collaborator's {@code tab} equal to the fixed lead tab label would shadow a lead sharing
* that space — refused outright.
*/
@Test
void aCollaboratorTabEqualToTheFixedLeadTabLabelRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("collaborator-is-lead.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
collaborators:
impostor:
tab: "lead"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("impostor"), "the message must name the offending collaborator");
assertTrue(e.getMessage().contains(FleetConfig.Leader.LEAD_TAB_LABEL),
"the message must name the fixed lead tab label it collides with");
}
// ── validatePanePlacementAgainstLeadTabs ────────────────────────────────────────────────────
/**
@@ -874,8 +923,12 @@ class FleetConfigTest {
assertTrue(e.getMessage().contains("gx10"), "the message must name the offending profile");
}
/**
* A lead is found by its fixed tab label regardless of its own {@code tab} field, so a
* {@code fleet.leaders} entry with no {@code tab} must still arm the guard.
*/
@Test
void aPanePlacedProfileWithNoLeadTabIsAllowed(@TempDir Path dir) throws Exception {
void aPanePlacedProfileWithNoLeadTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("pane-no-tab.yaml");
Files.writeString(f, """
bind:
@@ -888,9 +941,59 @@ class FleetConfigTest {
opus:
profile: gx10
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e = assertThrows(IllegalStateException.class,
cfg::validatePanePlacementAgainstLeadTabs);
assertTrue(e.getMessage().contains("gx10"), "the message must name the offending profile");
assertTrue(e.getMessage().contains("the only fix when a lead triggered this"),
"the message must say placement: tab is the only fix for a lead");
assertFalse(e.getMessage().contains("remove the tab from every fleet.leaders"),
"the message must not send the operator in a circle by advising a tab: removal");
}
/**
* {@code placement:} is optional, and {@link FleetConfig.Profile}'s own compact constructor
* defaults an absent or blank value to {@code "tab"}, so a profile naming no placement at all
* is tab-placed and the guard must not fire for it.
*/
@Test
void aProfileWithNoPlacementKeyDefaultsToTabPlacementAndIsAllowed(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("no-placement-key.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10: {}
fleet:
leaders:
opus:
profile: gx10
""");
FleetConfig cfg = FleetConfig.load(f);
assertTrue(cfg.profiles().get("gx10").tabPlacement(),
"Profile's compact constructor defaults an absent placement to \"tab\"");
assertDoesNotThrow(cfg::validatePanePlacementAgainstLeadTabs,
"a profile with no placement: key is tab-placed, not pane-placed");
}
/** Control: no {@code fleet.leaders} entry and no collaborator tab still starts fine. */
@Test
void aPanePlacedProfileWithAnEmptyFleetBlockIsAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("pane-empty-fleet.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
placement: pane
fleet: {}
""");
assertDoesNotThrow(() -> FleetConfig.load(f).validatePanePlacementAgainstLeadTabs(),
"a leader with no tab feeds nothing into the scanner, so pane placement is safe");
"no fleet.leaders entry and no collaborator tab means pane placement is safe");
}
@Test
@@ -24,6 +24,41 @@ public final class FakeHerdr implements HerdrClient {
/** The foreground PID of the one agent pane (term_a) in the canned {@code pane.process_info}. */
public static final long WORKER_PID = 4242;
/**
* A {@code detection} region of the current Claude Code TUI, whose input box is empty: a caret
* line between two rules, above the footer.
*/
public static final String IDLE_PROMPT_CARET = """
──────────────────────── lead: opus ─
❯
─────────────────────────────────────
lead: opus · Opus 5 (1M context) · ~/LTMS/claude-bridge
⏵⏵ auto mode on (shift+tab to cycle) · ← 1 agent""";
/** The same region with the operator's unsubmitted line still at the caret. */
public static final String DRAFTED_PROMPT_CARET = """
──────────────────────── lead: opus ─
❯ yes, send it to lead: opus
─────────────────────────────────────
lead: opus · Opus 5 (1M context) · ~/LTMS/claude-bridge
⏵⏵ auto mode on (shift+tab to cycle) · ← 1 agent""";
/** A {@code detection} region of an older TUI, which drew a bordered box, with that box empty. */
public static final String IDLE_PROMPT_BOX = """
⏺ done
╭────────────────────────╮
│ > │
╰────────────────────────╯
⏵⏵ auto mode on""";
/** The older TUI's bordered box, still holding the operator's unsubmitted line. */
public static final String DRAFTED_PROMPT_BOX = """
⏺ done
╭────────────────────────╮
│ > fix the issue when I │
╰────────────────────────╯
⏵⏵ auto mode on""";
private final ObjectMapper mapper = new ObjectMapper();
/**
* Thread-safe on purpose. Background loops — {@link dev.ltms.fleet.msg.ReplyPushLoop} and the
@@ -48,12 +83,15 @@ public final class FakeHerdr implements HerdrClient {
private final Map<String, String> processInfoErrorCodeFor = new ConcurrentHashMap<>();
private String tabCloseErrorCode = null;
private final Map<String, String> tabCloseErrorCodeFor = new ConcurrentHashMap<>();
private String workspaceListErrorCode = null;
private String agentSendErrorCode = null;
private boolean noPanes = false;
private volatile String agentStatus = "idle"; // steady-state agent.get status
private volatile String agentType = "claude"; // detected agent kind on agent.get; null = undetected
private volatile String agentSessionId = null; // agent_session.value on agent.get; null = omitted
private volatile String readText = "worker transcript tail"; // canned agent.read output
/** Canned {@code detection}-source output, or {@code null} to serve {@link #readText} there too. */
private volatile String detectionText = null;
private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane
private String pinnedStartTerminal;
private String pinnedStartPane;
@@ -142,6 +180,12 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/** Make {@code workspace.list} fail with this herdr error code; every other method still succeeds. */
public FakeHerdr workspaceListFailsWith(String code) {
this.workspaceListErrorCode = code;
return this;
}
/**
* Make {@code pane.list} report no panes at all — models a second herdr daemon (CB-185) that
* simply does not host the pane a {@link PaneLocator} is searching for.
@@ -200,12 +244,26 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/** The text {@code agent.read} returns (the CB-106 completion scrape). */
/**
* The text {@code agent.read} returns (the CB-106 completion scrape). It serves the
* {@code detection} source as well unless {@link #detectionText} overrides that one.
*/
public FakeHerdr readText(String text) {
this.readText = text;
return this;
}
/**
* Override the text {@code agent.read} returns for the {@code detection} source only — the
* prompt/footer tail herdr uses for status detection, a different region from the transcript the
* other sources carry. Needed by a test whose subject reads the input box, since one
* {@link #readText} cannot be both a worker's transcript and a lead's empty prompt.
*/
public FakeHerdr detectionText(String text) {
this.detectionText = text;
return this;
}
/** Make delivery ({@code agent.prompt} / {@code agent.send_keys}) fail with this error code. */
public FakeHerdr agentSendFailsWith(String code) {
this.agentSendErrorCode = code;
@@ -302,11 +360,17 @@ public final class FakeHerdr implements HerdrClient {
case "ping" -> mapper.readTree(
("{\"type\":\"pong\",\"version\":\"%s\",\"protocol\":%d}")
.formatted(pingVersion, pingProtocol));
case "workspace.list" -> mapper.readTree(("""
case "workspace.list" -> {
if (workspaceListErrorCode != null) {
throw new HerdrException("herdr error [" + workspaceListErrorCode + "]: workspace.list failed",
workspaceListErrorCode, null);
}
yield mapper.readTree(("""
{"type":"workspace_list","workspaces":[
{"workspace_id":"w1","label":"dev-mgnl","focused":true,"pane_count":7,"agent_status":"unknown"},
{"workspace_id":"w2","label":"ltms","focused":false,"pane_count":5,"agent_status":"done"}%s]}""")
.formatted(extraWorkspaces.isEmpty() ? "" : "," + String.join(",", extraWorkspaces)));
}
case "agent.list" -> mapper.readTree(("""
{"type":"agent_list","agents":[
{"terminal_id":"term_a","agent":"claude","agent_status":"idle",
@@ -347,8 +411,13 @@ public final class FakeHerdr implements HerdrClient {
"agent_status":"%s","workspace_id":"w2","tab_id":"w2:t7","pane_id":"w2:p7"%s}}""")
.formatted(agentField, agentStatus, sessionField));
}
case "agent.read" -> mapper.readTree(mapper.writeValueAsString(
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", readText))));
case "agent.read" -> {
Object source = params instanceof Map<?, ?> m ? m.get("source") : null;
String text = "detection".equals(source) && detectionText != null
? detectionText : readText;
yield mapper.readTree(mapper.writeValueAsString(
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", text))));
}
case "agent.start" -> {
if (onAgentStart != null) {
onAgentStart.run();
@@ -158,15 +158,24 @@ class LeadTabScannerTest {
return Map.of("lead: opus-5.0", "opus-5.0", "lead: gpt-sol-5.6", "gpt-sol-5.6");
}
/** {@code twoLeads()}'s own space — every lead-label fixture below lives here unless noted. */
private static final String MAIN_SPACE = "main";
/** Wraps a flat label → name map under one space, the shape {@link LeadTabScanner} now takes. */
private static Map<String, Map<String, String>> inSpace(String space, Map<String, String> labelToName) {
return Map.of(space, labelToName);
}
private LeadTabScanner scanner(TopologyHerdr herdr, Map<String, String> tabToName,
AtomicLong clock) {
return new LeadTabScanner(herdr, tabToName, Set.of("fleetd-workers"), TTL, clock::get);
return new LeadTabScanner(herdr, inSpace(MAIN_SPACE, tabToName), Set.of("fleetd-workers"),
TTL, clock::get);
}
private LeadTabScanner scannerWithCollaborators(TopologyHerdr herdr, Map<String, String> tabToName,
Map<String, String> collaboratorTabToName,
AtomicLong clock) {
return new LeadTabScanner(herdr, tabToName, collaboratorTabToName,
return new LeadTabScanner(herdr, inSpace(MAIN_SPACE, tabToName), collaboratorTabToName,
Set.of("fleetd-workers"), TTL, clock::get);
}
@@ -237,14 +246,61 @@ class LeadTabScannerTest {
.tab("w1:t2", "w1", "worker: gx10 #1")
.pane("w1:p1", "w1:t1", "term_opus")
.pane("w1:p2", "w1:t2", "term_worker");
LeadTabScanner s = new LeadTabScanner(herdr, Map.of("lead: opus-5.0", "opus-5.0"),
Set.of(), TTL, new AtomicLong()::get);
LeadTabScanner s = new LeadTabScanner(herdr,
inSpace("fleet", Map.of("lead: opus-5.0", "opus-5.0")), Set.of(), TTL,
new AtomicLong()::get);
assertEquals("opus-5.0", s.get().get("term_opus"),
"a lead sharing the members' workspace is still discovered — the label, not the "
+ "workspace, is what matches it");
}
// ── fleetd #770: space is the uniqueness boundary, not the label alone ──────────────────────
/**
* Two leads can share the exact same label (the fixed {@code lead} tab label) as long as they
* sit in different spaces — each tab resolves to its own space's lead, never the other one's.
*/
@Test
void aTabLabelledLeadResolvesToItsOwnSpacesLeadNotTheOtherSpaces() {
TopologyHerdr herdr = new TopologyHerdr()
.workspace("wa", "space-a")
.workspace("wb", "space-b")
.tab("wa:t1", "wa", "lead")
.tab("wb:t1", "wb", "lead")
.pane("wa:p1", "wa:t1", "term_a")
.pane("wb:p1", "wb:t1", "term_b");
Map<String, Map<String, String>> leadLabelsBySpace = Map.of(
"space-a", Map.of("lead", "alpha"),
"space-b", Map.of("lead", "beta"));
LeadTabScanner s = new LeadTabScanner(herdr, leadLabelsBySpace, Set.of(), TTL,
new AtomicLong()::get);
Map<String, String> leads = s.get();
assertEquals("alpha", leads.get("term_a"), "space-a's tab must resolve to space-a's lead");
assertEquals("beta", leads.get("term_b"), "space-b's tab must resolve to space-b's lead");
}
/**
* A lead's deprecated legacy {@code tab:} label is still matched, but only within that lead's
* own configured space — exactly the shape {@code FleetdAssembly} builds via {@code
* Leader.acceptedLabels()}.
*/
@Test
void aLegacyTabLabelStillResolvesWithinItsOwnSpace() {
TopologyHerdr herdr = new TopologyHerdr()
.workspace("w1", "fleet")
.tab("w1:t1", "w1", "lead: opus")
.pane("w1:p1", "w1:t1", "term_opus");
Map<String, Map<String, String>> leadLabelsBySpace =
Map.of("fleet", Map.of("lead", "opus", "lead: opus", "opus"));
LeadTabScanner s = new LeadTabScanner(herdr, leadLabelsBySpace, Set.of(), TTL,
new AtomicLong()::get);
assertEquals("opus", s.get().get("term_opus"),
"the deprecated tab label must still resolve this lead within its own space");
}
@Test
void aLabelWithNoConfiguredEntryIsIgnored() {
TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main")
@@ -0,0 +1,132 @@
package dev.ltms.fleet.herdr;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/** Reading a Claude Code input box, so a paste-and-submit delivery never submits the operator's draft. */
class PromptBoxTest {
private static final String EMPTY = FakeHerdr.IDLE_PROMPT_CARET;
private static final String DRAFTED = FakeHerdr.DRAFTED_PROMPT_CARET;
// --- pure classification -------------------------------------------------
@Test
void anEmptyCaretLineIsAnEmptyBox() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(EMPTY));
}
@Test
void aCaretLineHoldingTextIsADraftAndCountsItsCharacters() {
PromptBox.Reading reading = PromptBox.classify(DRAFTED);
assertEquals(PromptBox.State.DRAFT, reading.state());
assertEquals("yes,sendittolead:opus".length(), reading.characters(),
"padding does not count — only what the operator typed");
}
@Test
void aBorderedBoxIsReadToo() {
assertEquals(PromptBox.State.EMPTY, PromptBox.classify(FakeHerdr.IDLE_PROMPT_BOX).state(),
"an older TUI draws a bordered box, and its panes must still be readable");
assertEquals(PromptBox.State.DRAFT, PromptBox.classify(FakeHerdr.DRAFTED_PROMPT_BOX).state());
}
@Test
void aSingleTypedCharacterIsADraft() {
assertEquals(PromptBox.State.DRAFT, PromptBox.classify("❯ f").state());
assertEquals(PromptBox.State.DRAFT, PromptBox.classify("│ > f │").state());
}
@Test
void aCursorBlockInAnOtherwiseEmptyBoxIsEmpty() {
assertEquals(PromptBox.State.EMPTY, PromptBox.classify("❯ █").state(),
"a terminal capture may leave the cursor cell in an empty box");
assertEquals(PromptBox.State.EMPTY, PromptBox.classify("│ > █ │").state());
}
@Test
void theBoxLineIsReadToItsEndWhateverFollowsIt() {
assertEquals(PromptBox.State.DRAFT, PromptBox.classify("❯ half a line\n ⏵⏵ auto mode on").state());
assertEquals(PromptBox.State.EMPTY, PromptBox.classify("❯\n ⏵⏵ auto mode on").state());
}
@Test
void theLastBoxLineOnThePaneIsTheLiveOne() {
assertEquals(PromptBox.State.DRAFT,
PromptBox.classify("❯ an earlier prompt\n⏺ its answer\n❯ typing now").state(),
"the detection region carries scrollback, so earlier prompts sit above the live box");
assertEquals(PromptBox.State.EMPTY,
PromptBox.classify("❯ an earlier prompt\n⏺ its answer\n❯").state());
}
@Test
void aMarkerPartWayAlongALineIsNotABox() {
assertEquals(PromptBox.State.UNREADABLE, PromptBox.classify("⏺ type ❯ to get a prompt").state(),
"a caret the operator quoted is transcript text, not an input box");
assertEquals(PromptBox.State.EMPTY, PromptBox.classify("⏺ type ❯ to get a prompt\n❯").state(),
"and it must not shadow the real box further down");
}
@Test
void aPaneWithNoBoxIsUnreadable() {
assertEquals(PromptBox.State.UNREADABLE, PromptBox.classify("garbled ansi noise").state());
assertEquals(PromptBox.State.UNREADABLE, PromptBox.classify("").state());
assertEquals(PromptBox.State.UNREADABLE, PromptBox.classify(null).state());
}
@Test
void aGeneratingTurnIsUnreadableEvenWithAnEmptyBox() {
assertEquals(PromptBox.State.UNREADABLE,
PromptBox.classify(EMPTY + "\n ✳ Thinking… (12s · esc to interrupt)").state());
}
@Test
void aGeneratingMarkerInScrollbackAboveTheBoxDoesNotMakeThePaneUnreadable() {
assertEquals(PromptBox.State.EMPTY,
PromptBox.classify(" ✳ Thinking… (12s · esc to interrupt)\n⏺ done\n" + EMPTY).state(),
"that marker survives in scrollback, and holding on it would hold every delivery forever");
}
// --- the gate ------------------------------------------------------------
@Test
void anEmptyBoxClearsTheGateAndReadsTheDetectionRegion() {
FakeHerdr herdr = new FakeHerdr().detectionText(EMPTY);
assertTrue(new PromptBox(new AgentControl(herdr)).clearToSubmit("term_a"));
@SuppressWarnings("unchecked")
var params = (java.util.Map<String, Object>) herdr.lastCall("agent.read").params();
assertEquals("detection", params.get("source"),
"the input box is drawn in the detection region, not in transcript scrollback");
}
@Test
void aDraftedBoxHoldsTheGate() {
assertFalse(new PromptBox(new AgentControl(new FakeHerdr().detectionText(DRAFTED))).clearToSubmit("term_a"));
}
@Test
void anUnreadablePaneHoldsTheGate() {
assertFalse(new PromptBox(new AgentControl(new FakeHerdr().detectionText("garbled"))).clearToSubmit("term_a"));
}
@Test
void aFailedReadHoldsTheGate() {
assertFalse(new PromptBox(new AgentControl(new FakeHerdr().healthy(false))).clearToSubmit("term_a"),
"a pane this cannot read must never be pasted into");
}
@Test
void theGateClearsAgainOnceTheBoxEmpties() {
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
PromptBox box = new PromptBox(new AgentControl(herdr));
assertFalse(box.clearToSubmit("term_a"));
herdr.detectionText(EMPTY);
assertTrue(box.clearToSubmit("term_a"));
}
}
@@ -89,6 +89,11 @@ class BackendOutageFlowTest {
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
.put("terminal_id", "term_primary").put("agent_status", "idle"));
}
if ("agent.read".equals(method)) {
// The lead-nudge paths read the input box before pasting into it.
return MAPPER.createObjectNode().set("read",
MAPPER.createObjectNode().put("text", FakeHerdr.IDLE_PROMPT_CARET));
}
if ("agent.prompt".equals(method)) {
prompts.add(params);
sendLatch.countDown();
@@ -59,6 +59,17 @@ class InjectorTest {
assertEquals(List.of("hello"), sent());
}
@Test
void deliveringToAMemberReadsNoPane() {
// No human types into a spawned member's pane, so its delivery path must not pay for a
// prompt-box read the way a lead's nudge paths do.
injector.enqueue(T, "task", TestTurnTokens.inert(T));
injector.onStatus(T, AgentStatus.IDLE);
assertEquals(List.of("task"), sent());
assertFalse(herdr.called("agent.read"), "a member delivery must not read its pane");
}
@Test
void holdsDeliveryUntilTheWorkerIsAvailable() {
// CB-113: idle alone is not enough — hold until the worker's MCP is connected (ready).
@@ -114,13 +114,13 @@ class LeadLauncherTest {
"name is lead-<name>-<nonce>-<seq>: " + startedName(herdr));
}
/** The tab is labelled with the configured `tab:` so the scanner finds the lead on the next resolve. */
/** The tab is labelled with the fixed lead tab label so the scanner finds the lead on the next resolve. */
@Test
void labelsTheTabWithTheConfiguredTabValue() {
void labelsTheTabWithTheFixedLeadTabLabel() {
FakeHerdr herdr = new FakeHerdr();
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
assertEquals("lead: opus",
assertEquals("lead",
((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"));
}
@@ -148,6 +148,25 @@ class LeadLauncherTest {
assertFalse(herdr.called("tab.close"), "a labelled tab WITH a live agent must never be closed");
}
/**
* The tab a live lead actually sits in still carries its deprecated legacy {@code tab:} label,
* not the fixed {@code lead} tab label a freshly auto-launched instance would get. Counting must
* still recognise it as the live lead via {@link FleetConfig.Leader#acceptedLabels()}, or a
* daemon restart would read it as missing and launch a second orchestrator next to the first.
*/
@Test
void aLiveLeadInALegacyLabelledTabIsCountedSoNothingIsLaunched() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus")
.withAgent("lead-opus", "term_lead", "wL:p1", "wL:t1");
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"the legacy-labelled live lead must be counted — nothing may be launched");
assertFalse(herdr.called("agent.start"),
"a tab label fixed to a constant must not blind the count to a legacy-labelled lead");
}
/**
* The reason liveness is not "does the label exist". A tab left labelled by a session that has
* since died must not block the relaunch, or one crash disables auto-launch permanently.
@@ -263,7 +282,7 @@ class LeadLauncherTest {
assertFalse(herdr.called("tab.close"), "a tab running an agent again must never be closed");
assertFalse(herdr.called("agent.start"), "the lead is live again — nothing to relaunch");
assertEquals("wL:t1", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("tab_id"));
assertEquals("lead: opus", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"),
assertEquals("lead", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"),
"the pending-close flag must be cleared once the tab is confirmed live again");
}
@@ -285,7 +304,7 @@ class LeadLauncherTest {
@Test
void aHandOpenedLeadWithTheConfiguredTabLabelCountsAsLive() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wX", "main")
.withWorkspace("wX", "fleet")
.withTab("wX", "wX:t1", "lead: opus")
.withAgent("hand-opened", "term_hand", "wX:p1", "wX:t1");
@@ -597,7 +616,7 @@ class LeadLauncherTest {
"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. */
/** The new tab is labelled with the fixed lead tab label, and AFTER the start. */
@Test
void relaunchLabelsTheNewTabAfterStarting() {
FakeHerdr herdr = new FakeHerdr();
@@ -606,7 +625,7 @@ class LeadLauncherTest {
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started);
assertEquals("lead: opus", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"));
assertEquals("lead", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"));
int startIndex = indexOfLastCall(herdr, "agent.start");
int renameIndex = indexOfLastCall(herdr, "tab.rename");
@@ -16,6 +16,7 @@ import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -1127,8 +1128,13 @@ class LeadRolloverTest {
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
// Only LEAD resolves to a configured name — every other terminal (e.g. the distinct
// term_cap_N terminals evictionCountsInProgressEntriesTowardTheCap opens) falls back to
// keying its own single-flight claim on the terminal itself, exactly like a terminal
// the live roster does not recognise.
return new LeadRollover(agents, spaces, launcher, () -> config, _ -> null,
_ -> LEAD_NAME, () -> liveLeadTerminals, nowMillis, () -> { }, runner);
t -> LEAD.equals(t) ? LEAD_NAME : null, () -> liveLeadTerminals, nowMillis, () -> { },
runner);
}
@Test
@@ -1317,7 +1323,7 @@ class LeadRolloverTest {
+ "so an operator reading status() has something to act on: " + status.detail());
}
// ---- fleetd #726 unit 3: confirm() single-flights per lead terminal -----------------------
// ---- fleetd #726 unit 3: confirm() single-flights one roll at a time per lead ---------------
@Test
@DisplayName("[SINGLE-FLIGHT 1] a second confirm() for the SAME lead terminal is refused with "
@@ -1343,8 +1349,8 @@ class LeadRolloverTest {
assertFalse(secondDecision.accepted());
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
assertTrue(secondDecision.detail().contains(LEAD), "the refusal detail must name the lead "
+ "terminal: " + secondDecision.detail());
assertTrue(secondDecision.detail().contains(LEAD_NAME), "the refusal detail must name "
+ "the single-flight key (the lead's name): " + secondDecision.detail());
assertTrue(secondDecision.detail().contains(first.token()), "the refusal detail must name "
+ "the token holding the claim: " + secondDecision.detail());
assertEquals(1, runner.heldCount(), "the refused confirm() must never have reached "
@@ -1502,4 +1508,160 @@ class LeadRolloverTest {
+ "this lead terminal could never be rolled again: " + retryDecision.reason() + " / "
+ retryDecision.detail());
}
// ---- fleetd #737 unit 4: confirm() single-flights per LEAD NAME, not per lead terminal ------
@Test
@DisplayName("[NAME-KEYED 1] a second confirm() for the SAME lead is refused with "
+ "ROLL_ALREADY_RUNNING even when it is opened from a DIFFERENT terminal, while the "
+ "first roll's continuation is still in flight")
void secondConfirmForTheSameLeadIsRefusedEvenFromADifferentTerminal() throws IOException {
FakeHerdr herdr = herdrReadyForAFullRoll(); // the held roll WOULD complete once run
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner();
Map<String, String> liveLeadTerminals = Map.of(NEW_TERMINAL, LEAD_NAME);
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
// Both LEAD and OTHER_LEAD resolve to the SAME configured lead name — modelling a roll
// that has already replaced the lead's pane: the fresh pane (here, OTHER_LEAD) is still
// the SAME lead, just a different terminal id.
Function<String, String> leadNameForTerminal = _ -> LEAD_NAME;
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, leadNameForTerminal, () -> liveLeadTerminals, fixedClock(clock), () -> { },
runner);
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision firstDecision = rollover.confirm(LEAD, first.token(), true);
assertTrue(firstDecision.accepted(), "expected approval; got: " + firstDecision.reason()
+ " / " + firstDecision.detail());
assertEquals(1, runner.heldCount(), "sanity: the first roll is held, not run yet");
LeadRollover.PendingRollover second = rollover.open(OTHER_LEAD,
"a second request for the SAME lead, opened from a DIFFERENT terminal");
LeadRollover.RollDecision secondDecision = rollover.confirm(OTHER_LEAD, second.token(), true);
assertFalse(secondDecision.accepted(), "a different terminal resolving to the SAME lead "
+ "name must still be refused — the single-flight claim is keyed on the lead's "
+ "name, not its terminal");
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
assertTrue(secondDecision.detail().contains(LEAD_NAME), "the refusal detail must name "
+ "the single-flight key (the lead's name): " + secondDecision.detail());
assertTrue(secondDecision.detail().contains(first.token()), "the refusal detail must "
+ "name the token holding the claim: " + secondDecision.detail());
assertEquals(1, runner.heldCount(), "the refused confirm() must never have reached "
+ "continuationRunner — it ran exactly once, for the first roll only");
runner.runNext(); // let the first (and only) held roll finish
// The claim must have been released once the first roll's continuation finished — a
// fresh request for the SAME lead, even opened from yet another terminal, is now
// approved.
LeadRollover.PendingRollover third = rollover.open(OTHER_LEAD, "retry after the first roll finished");
LeadRollover.RollDecision thirdDecision = rollover.confirm(OTHER_LEAD, third.token(), true);
assertTrue(thirdDecision.accepted(), "the claim must have been released once the first "
+ "roll finished: " + thirdDecision.reason() + " / " + thirdDecision.detail());
}
@Test
@DisplayName("[NAME-KEYED 2] the claim is released under the key it was taken under, even "
+ "when the old terminal has already dropped out of the live-lead roster by release "
+ "time — success (non-throwing) path")
void claimReleasedUnderCarriedKeyEvenWhenOldTerminalIsNoLongerInTheRosterSuccessPath() throws IOException {
// herdrReadyForAFullRoll() lets the RETRY below run its continuation to a genuine ROLLED
// completion (pane gone on the first check, a pinned relaunch) — needed because, unlike
// the other NAME-KEYED tests, this one's retry resolves a REAL name ("opus") and must not
// spin forever against this test's fixed, non-advancing clock (see this class's javadoc).
FakeHerdr herdr = herdrReadyForAFullRoll();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
// A config that CAN relaunch 'opus' — unlike emptyFleetConfig(), its fleet() is non-null,
// so resolveLaunchable(null) below returns null cleanly instead of throwing a NullPointer
// out of cfg.fleet() itself; this test needs the NON-throwing relaunch-refused exit, not
// an incidental NPE.
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
Map<String, String> roster = new HashMap<>();
roster.put(LEAD, LEAD_NAME);
Map<String, String> liveLeadTerminals = Map.of(NEW_TERMINAL, LEAD_NAME);
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, roster::get, () -> liveLeadTerminals, fixedClock(clock), () -> { }, Runnable::run);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
// The live roster drops the OLD terminal before this roll's continuation runs — the
// hazard fleetd #737 names: leadNameForTerminal reads the LIVE roster, and (with the
// synchronous runner this test injects) the continuation runs INSIDE this confirm()
// call, strictly after open() already resolved and carried the key.
roster.remove(LEAD);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.RELAUNCH_FAILED, status.state(), "sanity: leadName "
+ "resolved to null (the roster had already dropped the old terminal), so "
+ "relaunch(null) fails cleanly, through the non-throwing exit this test means "
+ "to cover: " + status.detail());
// A fresh terminal for the SAME lead name, confirmed live again — proves the claim this
// roll took under the lead's name was actually released, not left stuck under whatever
// a release that RECOMPUTES the key via leadNameForTerminal (now null for the old
// terminal) would have tried to remove instead.
roster.put(OTHER_LEAD, LEAD_NAME);
LeadRollover.PendingRollover retry = rollover.open(OTHER_LEAD, "a fresh pane for the same lead");
LeadRollover.RollDecision retryDecision = rollover.confirm(OTHER_LEAD, retry.token(), true);
assertTrue(retryDecision.accepted(), "the claim taken under the lead's name must have "
+ "been released even though the OLD terminal no longer resolved to any name at "
+ "release time: " + retryDecision.reason() + " / " + retryDecision.detail());
LeadRollover.RollStatus retryStatus = rollover.status(retry.token());
assertEquals(LeadRollover.RollState.ROLLED, retryStatus.state(), "sanity: this retry's "
+ "own claim (also 'opus') must not have been blocked by a leftover claim from "
+ "the first roll: " + retryStatus.detail());
}
@Test
@DisplayName("[NAME-KEYED 3] the claim is released under the key it was taken under on the "
+ "thrown-exception path too, even when the old terminal has already dropped out of "
+ "the live-lead roster")
void claimReleasedUnderCarriedKeyEvenWhenOldTerminalIsNoLongerInTheRosterThrowingPath() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
fake.paneCloseFailsWith("permission_denied"); // a real failure, not an already-gone code
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
AgentControl agents = new AgentControl(fake);
WorkspaceControl spaces = new WorkspaceControl(fake);
LeadLauncher launcher = fakeLauncher(fake, emptyFleetConfig());
Map<String, String> roster = new HashMap<>();
roster.put(LEAD, LEAD_NAME);
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, roster::get, Map::of, fixedClock(clock), () -> { }, Runnable::run);
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
// Same hazard as [NAME-KEYED 2] above, now exercised on the thrown-exception exit.
roster.remove(LEAD);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, first.token(), true);
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
+ "inside the deferred continuation, which this test's synchronous runner has "
+ "already run to completion by the time confirm() returns");
LeadRollover.RollStatus status = rollover.status(first.token());
assertEquals(LeadRollover.RollState.FAILED, status.state(), "sanity: the continuation "
+ "threw and left a terminal FAILED outcome: " + status.detail());
// A fresh terminal for the SAME lead name, confirmed live again — proves the claim this
// roll took under the lead's name was actually released, not left stuck under whatever
// a release that RECOMPUTES the key via leadNameForTerminal (now null for the old
// terminal) would have tried to remove instead.
roster.put(OTHER_LEAD, LEAD_NAME);
LeadRollover.PendingRollover retry = rollover.open(OTHER_LEAD, "retry after the throw");
LeadRollover.RollDecision retryDecision = rollover.confirm(OTHER_LEAD, retry.token(), true);
assertTrue(retryDecision.accepted(), "the claim must have been released even though the "
+ "continuation threw AND the old terminal no longer resolved to any name at "
+ "release time: " + retryDecision.reason() + " / " + retryDecision.detail());
}
}
@@ -17,6 +17,7 @@ import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
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;
@@ -269,6 +270,60 @@ class FleetMcpAuthzTest {
"a spawned member's own terminal must stay unreachable, even once the classifier is real");
}
// --- fleetd #743: the observer SEND matrix, over MCP's denyFor -------------------------------
private static final Principal OBSERVER = Principal.observer("term_observer", 700);
/**
* Wires one real {@link CallerResolver} that recognises a lead, a collaborator, and a live
* spawned worker, leaving "term_other_observer" classified as none of them — so the same
* wiring both denies an observer's {@code SEND} to every privileged role and grants it to
* another unclassified pane, proving the refusals are the rule and not a missing fixture.
*/
@Test
void anObserverMaySendOnlyToAnotherObserverNeverToALeadWorkerOrArchitect() {
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> Map.of("term_lead_known", "lead-x"), new MemberRegistry(null),
t -> "term_a".equals(t) ? MemberRole.DEV : null,
() -> Map.of("term_collab_known", "ops2"));
FleetMcp m = mcp(true, callers);
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_lead_known"),
"an observer must never reach a lead's terminal");
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_collab_known"),
"an observer must never reach a collaborator's terminal");
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_a"),
"an observer must never reach a live spawned member's terminal");
// CONTROL: the same wiring, the same denyFor call, a target recognised as none of the
// three privileged roles above -- this is what proves the three refusals above are the
// rule working, not a classifier that refuses every target regardless of what it is.
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_other_observer"),
"an observer must reach another pane that resolves as an observer itself");
}
/**
* As {@link #anObserverMaySendOnlyToAnotherObserverNeverToALeadWorkerOrArchitect}, for a
* terminal bound to a configured architect slot but hosting no live spawned-member session --
* the case {@link CallerResolver#resolve} itself treats separately from a live worker/architect.
*/
@Test
void anObserverMayNotSendToABoundArchitectSlotEither() {
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(Map.of(),
Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null));
assertTrue(members.bind("architect:lead-designer", "term_bound_architect"));
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, Map::of,
members, t -> null, Map::of);
FleetMcp m = mcp(true, callers);
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_bound_architect"),
"a terminal bound to a configured architect slot must stay unreachable to an observer");
// CONTROL: the same wiring, a target the bind above never touched.
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_other_observer"));
}
@Test
void theLegacyConstructorLeavesTheGateOpen() {
// The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced.
@@ -453,6 +508,28 @@ class FleetMcpAuthzTest {
assertFalse(FleetMcp.membersVisibleTo(ANON), "authenticated as nothing must not see it either");
}
// --- who may see fleet_list's panes array ----------------------------------------------------
/**
* {@link FleetMcp#panesVisibleTo} is the whole policy decision for {@code fleet_list}'s
* {@code panes} array: visible to every role that may {@code SEND} to some other pane -- the
* primary, an architect, a collaborator, and an observer (to another observer pane only, with
* its row filtered and reduced -- see {@code listFleet}) -- never a worker, never an
* anonymous caller.
*/
@Test
void onlyPrimaryArchitectCollaboratorAndObserverMaySeeThePanesArray() {
assertTrue(FleetMcp.panesVisibleTo(PRIMARY), "the primary must see the panes array");
assertTrue(FleetMcp.panesVisibleTo(ARCH_DESIGN), "an architect must see the panes array");
assertTrue(FleetMcp.panesVisibleTo(COLLABORATOR), "a collaborator must see its own peer roster");
assertFalse(FleetMcp.panesVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the panes array");
assertTrue(FleetMcp.panesVisibleTo(Principal.observer("term_obs", 700)),
"an observer holds SEND to another observer pane, so it must see the (filtered, "
+ "reduced) panes array");
assertFalse(FleetMcp.panesVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* Same reasoning as {@link #theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}: the
* predicate above can be perfectly correct while the one production call site never asks it.
@@ -56,9 +56,14 @@ class FleetMcpHandoverTest {
private static final String LEAD = "term_lead";
private static final String OTHER_LEAD = "term_other_lead";
private static final String LEAD_OWNER = "leader:lead";
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
/** Fed to every direct {@code FleetMcp.handover} call below — none of this class's own tests
* exercise ticket/ask ownership, so a single instance with no delegations is enough. */
private final MessageService messages = new MessageService(agents, new Injector(agents),
new Rendezvous(), new InMemoryReplyInbox());
private FleetMcp mcp;
@AfterEach
@@ -151,16 +156,16 @@ class FleetMcpHandoverTest {
@DisplayName("with leadRollover: absent, every action returns a clean NOT_CONFIGURED refusal and never throws")
void nullLeadRolloverRefusesCleanlyForEveryAction() {
McpSchema.CallToolResult open = assertDoesNotThrow(
() -> FleetMcp.handover(null, LEAD, Map.of("action", "open")));
() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of("action", "open")));
assertFalse(open.isError(), "a refusal is not a protocol error: " + textOf(open));
assertTrue(textOf(open).contains("NOT_CONFIGURED"), textOf(open));
McpSchema.CallToolResult confirm = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
McpSchema.CallToolResult confirm = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
Map.of("action", "confirm", "token", "whatever")));
assertFalse(confirm.isError());
assertTrue(textOf(confirm).contains("NOT_CONFIGURED"), textOf(confirm));
McpSchema.CallToolResult cancel = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
McpSchema.CallToolResult cancel = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
Map.of("action", "cancel", "token", "whatever")));
assertFalse(cancel.isError());
assertTrue(textOf(cancel).contains("NOT_CONFIGURED"), textOf(cancel));
@@ -169,11 +174,11 @@ class FleetMcpHandoverTest {
@Test
@DisplayName("a blank/unknown action is a clean tool error, never an exception")
void unknownActionIsACleanError() {
McpSchema.CallToolResult missing = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD, Map.of()));
McpSchema.CallToolResult missing = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of()));
assertTrue(missing.isError());
McpSchema.CallToolResult bogus = assertDoesNotThrow(
() -> FleetMcp.handover(null, LEAD, Map.of("action", "bogus")));
() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of("action", "bogus")));
assertTrue(bogus.isError());
}
@@ -229,7 +234,7 @@ class FleetMcpHandoverTest {
Files.writeString(handover, "not written yet");
LeadRollover rollover = newRollover(handover.toString());
McpSchema.CallToolResult openResult = FleetMcp.handover(rollover, LEAD, Map.of("action", "open"));
McpSchema.CallToolResult openResult = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"));
assertFalse(openResult.isError(), textOf(openResult));
String token = extractToken(textOf(openResult));
@@ -238,13 +243,13 @@ class FleetMcpHandoverTest {
Thread.sleep(50);
Files.writeString(handover, "the real handover content");
McpSchema.CallToolResult wrongCaller = FleetMcp.handover(rollover, OTHER_LEAD,
McpSchema.CallToolResult wrongCaller = FleetMcp.handover(rollover, messages, OTHER_LEAD, "leader:other-lead",
Map.of("action", "confirm", "token", token));
assertFalse(wrongCaller.isError(), "a refusal is a legitimate outcome, not a protocol error");
assertTrue(textOf(wrongCaller).contains("NOT_YOUR_ROLLOVER"),
"a different lead terminal confirming must surface NOT_YOUR_ROLLOVER: " + textOf(wrongCaller));
McpSchema.CallToolResult confirmed = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult confirmed = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "confirm", "token", token));
assertFalse(confirmed.isError(), textOf(confirmed));
assertTrue(textOf(confirmed).contains("\"accepted\":true"),
@@ -258,7 +263,7 @@ class FleetMcpHandoverTest {
void cancelUnknownTokenIsCleanNotAFailure() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "cancel", "token", "does-not-exist"));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"cancelled\":false"), textOf(r));
@@ -269,9 +274,9 @@ class FleetMcpHandoverTest {
@DisplayName("cancel on a token actually opened reports cancelled:true")
void cancelKnownTokenSucceeds() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
String token = extractToken(textOf(FleetMcp.handover(rollover, LEAD, Map.of("action", "open"))));
String token = extractToken(textOf(FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"))));
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "cancel", "token", token));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"cancelled\":true"), textOf(r));
@@ -282,7 +287,7 @@ class FleetMcpHandoverTest {
@Test
@DisplayName("status on a null LeadRollover is a clean NOT_CONFIGURED refusal, never a throw")
void statusWithNullLeadRolloverRefusesCleanly() {
McpSchema.CallToolResult r = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
McpSchema.CallToolResult r = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
Map.of("action", "status", "token", "whatever")));
assertFalse(r.isError());
assertTrue(textOf(r).contains("NOT_CONFIGURED"), textOf(r));
@@ -293,7 +298,7 @@ class FleetMcpHandoverTest {
void statusOnUnknownTokenReportsUnknown() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "status", "token", "does-not-exist"));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"state\":\"UNKNOWN\""), textOf(r));
@@ -303,9 +308,9 @@ class FleetMcpHandoverTest {
@DisplayName("status on a token that is still pending (opened, not confirmed) reports PENDING")
void statusOnPendingTokenReportsPending() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
String token = extractToken(textOf(FleetMcp.handover(rollover, LEAD, Map.of("action", "open"))));
String token = extractToken(textOf(FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"))));
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "status", "token", token));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"state\":\"PENDING\""), textOf(r));
@@ -323,5 +328,40 @@ class FleetMcpHandoverTest {
"the tool's own description must advertise the 'status' action: " + tool.description());
}
// --- unit 5: open() reports outstanding tickets and open asks ------------------------------
@Test
@DisplayName("open() keeps token/handoverPath/requestedAtMillis and reports empty outstanding collections when the caller has nothing")
void openReportsEmptyOutstandingCollectionsWhenCallerHasNothing() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "open"));
assertFalse(r.isError(), textOf(r));
String json = textOf(r);
assertTrue(json.contains("\"token\":"), json);
assertTrue(json.contains("\"handoverPath\":"), json);
assertTrue(json.contains("\"requestedAtMillis\":"), json);
assertTrue(json.contains("\"outstandingTickets\":[]"),
"a caller with nothing gets an empty array, not an absent key: " + json);
assertTrue(json.contains("\"openAsks\":[]"),
"a caller with nothing gets an empty array, not an absent key: " + json);
}
@Test
@DisplayName("open() reports an owned pending ticket with its phase and target")
void openReportsAnOwnedPendingTicketWithItsPhase() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
String ticket = messages.sendAsync("term_worker", "a task", null, Principal.leader("lead", LEAD, 1));
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "open"));
assertFalse(r.isError(), textOf(r));
String json = textOf(r);
assertTrue(json.contains("\"ticket\":\"" + ticket + "\""), json);
assertTrue(json.contains("\"phase\":\"PENDING\""), json);
assertTrue(json.contains("\"target\":\"term_worker\""), json);
}
// --- acceptance 7 (wiring) is covered by FleetdLeadRolloverWiringTest, unchanged -----------
}
@@ -0,0 +1,142 @@
package dev.ltms.fleet.mcp;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
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.member.ClaudeCodeLauncher;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.session.FakeWorktrees;
import dev.ltms.fleet.session.SessionManager;
import io.modelcontextprotocol.client.McpClient;
import io.modelcontextprotocol.client.McpSyncClient;
import io.modelcontextprotocol.client.transport.HttpClientStreamableHttpTransport;
import io.modelcontextprotocol.spec.McpClientTransport;
import io.modelcontextprotocol.spec.McpSchema;
import org.eclipse.jetty.server.Server;
import org.eclipse.jetty.server.ServerConnector;
import org.eclipse.jetty.servlet.ServletContextHandler;
import org.eclipse.jetty.servlet.ServletHolder;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #743, driven end to end: a real MCP client over a real HTTP transport, resolved by the
* real {@link CallerResolver} to {@link dev.ltms.fleet.auth.Role#OBSERVER}, sending to another
* unclassified pane. {@link FleetMcpAuthzTest} already proves {@code denyFor} grants this case and
* that the handler calls {@code attributeIfObserver}; this test is the one path that also proves
* the grant is not dead at a second gate (CB-548's two-gate trap) by driving the real
* {@link Injector} to the point of its real herdr call, and reads the exact text the receiving
* pane would see.
*/
class FleetMcpObserverSendDeliveryTest {
private static final String TARGET = "term_other_observer";
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
private final Injector injector = new Injector(agents);
private final Rendezvous rendezvous = new Rendezvous();
private final MessageService messages = new MessageService(agents, injector, rendezvous,
new InMemoryReplyInbox());
private FleetMcp mcp;
private Server server;
@AfterEach
void tearDown() throws Exception {
if (server != null) {
server.stop();
}
if (mcp != null) {
mcp.close();
}
}
@Test
void anObserversSendIsAttributedAndReachesTheRealInjector() throws Exception {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(agents, new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> "tok");
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
// The fake's pane list carries a second pane, "term_shell", whose shell pid is 9001 and
// which hosts no agent -- a herdr-owned pane recognised as no lead, architect, collaborator
// or live spawned member, so the real resolver lands it on the observer floor.
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 9001L);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null));
mcp = new FleetMcp(messages, workers, sessions, identity, sessions.asPresence(),
new PrimaryRegistry(null), callers, FleetMcp.AuthorizationMode.ENFORCED,
null, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), null, FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), List.of(), null);
ServletContextHandler handler = new ServletContextHandler();
handler.setContextPath("/");
handler.addServlet(new ServletHolder(mcp.servlet()), "/mcp");
server = new Server(0);
server.setHandler(handler);
server.start();
String baseUrl = "http://127.0.0.1:"
+ ((ServerConnector) server.getConnectors()[0]).getLocalPort();
McpSchema.CallToolResult result = sendFleetSend(baseUrl, TARGET, "hi there");
assertFalse(result.isError(), "an observer sending to another observer must be accepted: "
+ textOf(result));
long waiterDeadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(TARGET) && System.currentTimeMillis() < waiterDeadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(TARGET), "the async send must have opened its rendezvous waiter");
injector.onStatus(TARGET, AgentStatus.IDLE); // drives the real delivery attempt to herdr
long deliveryDeadline = System.currentTimeMillis() + 3000;
while (!herdr.called("agent.prompt") && System.currentTimeMillis() < deliveryDeadline) {
Thread.sleep(5);
}
assertTrue(herdr.called("agent.prompt"), "the delivery attempt must have reached herdr");
@SuppressWarnings("unchecked")
Map<String, Object> params = (Map<String, Object>) herdr.lastCall("agent.prompt").params();
assertEquals("[fleet_send from observer term_shell]\nhi there", params.get("text"),
"the receiving pane must see the sender's own daemon-resolved terminal, never a raw "
+ "echo of the content and never a client-supplied name");
}
private static McpSchema.CallToolResult sendFleetSend(String baseUrl, String target, String content) {
McpClientTransport transport = HttpClientStreamableHttpTransport.builder(baseUrl)
.endpoint("/mcp")
.build();
try (McpSyncClient client = McpClient.sync(transport).build()) {
client.initialize();
return client.callTool(McpSchema.CallToolRequest.builder("fleet_send")
.arguments(Map.of("sessionId", target, "content", content, "wait", false))
.build());
}
}
private static String textOf(McpSchema.CallToolResult r) {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
}
@@ -1,5 +1,6 @@
package dev.ltms.fleet.mcp;
import dev.ltms.fleet.Fleetd;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.Principal;
@@ -1110,6 +1111,253 @@ class FleetMcpTest {
assertTrue(asArchitect.contains("\"sessionId\":\"term_collab\""), asArchitect);
}
// --- fleet_list's panes array -----------------------------------------------------------------
/** Calls the canonical {@code listFleet} overload directly, so a test can set the pane-discovery
* payload and its visibility independently of a real {@code Principal} / MCP exchange. */
private static McpSchema.CallToolResult listFleetWithPanes(FakeHerdr h, FleetMcp.PaneSource panes,
boolean panesVisible) {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
return FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
Map.of(), "", Map.of(), false,
FleetMcp.CoordinationSource.none(), false, true, true, panes, panesVisible);
}
/**
* A pane whose tab herdr reports with a label gets that label and the exact terminal id
* {@code fleet_send} takes as a target, carried as {@code sessionId}. fleetd #771: the pane's
* workspace carries its herdr space name too, next to {@code workspaceId}.
*/
@Test
void listReportsAPaneRowWithItsTabLabelAndSendableSessionId() {
FakeHerdr h = new FakeHerdr().withTab("w2", "w2:t7", "trinotes");
MemberPresence presence = new MemberPresence();
presence.markPresent("term_a");
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
() -> new PaneLocator(h).tabLabelsByTabId(),
() -> new PaneLocator(h).workspaceLabelsByWorkspaceId(),
Fleetd.deliverableTo(presence, Map::of, Map::of), _ -> false, _ -> false);
String out = textOf(listFleetWithPanes(h, panes, true));
assertTrue(out.contains("\"panes\":["), out);
assertTrue(out.contains("\"sessionId\":\"term_a\""), out);
assertTrue(out.contains("\"label\":\"trinotes\""), out);
assertTrue(out.contains("\"workspaceLabel\":\"ltms\""),
"term_a's agent lives on workspace w2, whose herdr label is \"ltms\": " + out);
assertTrue(out.contains("\"deliverable\":true"), out);
}
/**
* A pane whose tab carries no label known to herdr still gets a row -- a missing label must
* never throw, and must never drop the pane from the array, only report a {@code null} label.
* Pairs with a {@code deliverable} false reading when the target is neither present, a lead,
* nor a collaborator.
*/
@Test
void listReportsAPaneRowWithANullLabelWhenHerdrHasNoneAndNotDeliverable() {
FakeHerdr h = new FakeHerdr();
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
() -> new PaneLocator(h).tabLabelsByTabId(),
() -> new PaneLocator(h).workspaceLabelsByWorkspaceId(),
Fleetd.deliverableTo(new MemberPresence(), Map::of, Map::of), _ -> false, _ -> false);
String out = textOf(listFleetWithPanes(h, panes, true));
assertTrue(out.contains("\"panes\":["), out);
assertTrue(out.contains("\"sessionId\":\"term_a\""), out);
assertTrue(out.contains("\"label\":null"), out);
assertTrue(out.contains("\"deliverable\":false"), out);
}
/**
* fleetd #771: a pane whose agent lives in a workspace that {@code workspace.list} does not
* report (an unknown {@code workspaceId}) still gets a row -- the lookup miss must never throw,
* and must never drop the pane, only report a {@code null} "workspaceLabel".
*/
@Test
void listReportsANullWorkspaceLabelForAnUnknownWorkspaceId() {
// withAgent seeds its pane under workspace_id "wQ", which the fake's workspace.list never
// reports (only "w1"/"w2") -- modelling a workspace the lookup has no entry for.
FakeHerdr h = new FakeHerdr().withAgent("claude-x", "term_unknown_ws", "wQ:p1", "wQ:t1");
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
() -> new PaneLocator(h).tabLabelsByTabId(),
() -> new PaneLocator(h).workspaceLabelsByWorkspaceId(),
_ -> false, _ -> false, _ -> false);
String out = textOf(listFleetWithPanes(h, panes, true));
assertTrue(out.contains("\"sessionId\":\"term_unknown_ws\""), out);
assertTrue(out.contains("\"workspaceId\":\"wQ\""), out);
assertTrue(out.contains("\"workspaceLabel\":null"),
"an unknown workspaceId must project a null workspaceLabel, not throw or drop the row: " + out);
}
/** A caller this role may not show the array to gets no {@code panes} key at all. */
@Test
void listOmitsThePanesArrayWhenTheCallerMayNotSeeIt() {
FakeHerdr h = new FakeHerdr();
String out = textOf(listFleetWithPanes(h, FleetMcp.PaneSource.none(), false));
assertFalse(out.contains("\"panes\""), out);
}
/**
* The tab-label scan behind {@code panes} shares no failure path with the rest of
* {@code listFleet} -- a {@code workspace.list}/{@code tab.list} failure costs only the
* labels in the {@code panes} row (each renders {@code null}), never the {@code leads}/
* {@code members} arrays, which never needed that scan at all. fleetd #771: the same
* {@code workspace.list} failure costs {@code workspaceLabel} the same way.
*/
@Test
void listStillReportsEveryOtherArrayWhenTheLabelScanFails() {
FakeHerdr h = new FakeHerdr().workspaceListFailsWith("unavailable");
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
() -> new PaneLocator(h).tabLabelsByTabId(),
() -> new PaneLocator(h).workspaceLabelsByWorkspaceId(),
Fleetd.deliverableTo(new MemberPresence(), Map::of, Map::of), _ -> false, _ -> false);
McpSchema.CallToolResult res = listFleetWithPanes(h, panes, true);
assertNotEquals(Boolean.TRUE, res.isError(), textOf(res));
String out = textOf(res);
assertTrue(out.contains("\"panes\":["), out);
assertTrue(out.contains("\"sessionId\":\"term_a\""), out);
assertTrue(out.contains("\"label\":null"), out);
assertTrue(out.contains("\"workspaceLabel\":null"),
"a workspace.list failure must not cost the leads/members arrays, only a null "
+ "workspaceLabel: " + out);
assertTrue(out.contains("\"leads\":[]"), "a label-scan failure must not cost the leads array: " + out);
assertTrue(out.contains("\"members\":[]"), "a label-scan failure must not cost the members array: " + out);
}
/**
* fleetd #756: a pane bound to a configured architect slot with no live member session must
* report {@code role: "architect"}, read from {@link CallerResolver#boundToArchitectSlot} —
* the same classifier {@link CallerResolver#sendableObserverTarget} refuses as a {@code SEND}
* target — rather than falling through to {@code "observer"}.
*/
@Test
void listReportsArchitectForASlotBoundPaneWithNoLiveMember() {
FakeHerdr h = new FakeHerdr()
.withAgent("claude-arch", "term_unoccupied_architect", "w2:pArch", "w2:tArch");
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
Map::of, Map::of, _ -> false, "term_unoccupied_architect"::equals, _ -> false);
String out = textOf(listFleetWithPanes(h, panes, true));
assertTrue(out.contains("\"sessionId\":\"term_unoccupied_architect\""), out);
assertTrue(out.contains("\"role\":\"architect\""),
"a slot-bound pane with no live session must read \"architect\", not the generic "
+ "\"observer\" fallback: " + out);
}
/**
* Calls the canonical {@code listFleet} overload directly with an explicit {@code leads}/
* {@code collaborators} payload and {@code callerIsObserver}, mirroring exactly what the real
* {@code fleet_list} handler computes for an observer caller: {@code panesVisible} true,
* {@code leadsVisible}/{@code membersVisible}/{@code collaboratorsVisible} false.
*/
private static McpSchema.CallToolResult listFleetAsObserver(FakeHerdr h, Map<String, String> leads,
Map<String, String> collaborators, FleetMcp.PaneSource panes) {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
return FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
leads, "", collaborators, false,
FleetMcp.CoordinationSource.none(), false, false, false, panes, true, true);
}
/**
* fleetd #758: an observer's {@code fleet_list} now carries a {@code panes} key, filtered to
* {@link CallerResolver#sendableObserverTarget} (so a lead's pane, a spawned member's pane, a
* collaborator's pane, and an unoccupied architect-slot pane are all absent) and every
* surviving row reduced to exactly {@code sessionId}, {@code label}, {@code status},
* {@code role}, {@code deliverable} — never {@code paneId}, {@code workspaceId},
* {@code workspaceLabel} (fleetd #771 — a space name is host shape, a stronger disclosure than
* a pane id, so it stays out of the reduced row too), {@code tabId}, {@code agentType}, or
* {@code cwd}.
*/
@Test
void listFiltersAndReducesThePanesArrayForAnObserver() {
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(Map.of(),
Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null));
assertTrue(members.bind("architect:lead-designer", "term_architect_pane"));
CallerResolver callers = CallerResolver.withLeadsAndMembers(null, false, null,
() -> Map.of("term_lead_pane", "fleet01-lead"), members,
t -> "term_member_pane".equals(t) ? MemberRole.DEV : null,
() -> Map.of("term_collab_pane", "ops"));
FakeHerdr h = new FakeHerdr()
.withAgent("claude-sendable", "term_sendable", "w2:pS", "w2:tS")
.withAgent("claude-lead", "term_lead_pane", "w2:pL", "w2:tL")
.withAgent("claude-member", "term_member_pane", "w2:pM", "w2:tM")
.withAgent("claude-collab", "term_collab_pane", "w2:pC", "w2:tC")
.withAgent("claude-arch", "term_architect_pane", "w2:pA", "w2:tA")
.withTab("w2", "w2:tS", "trinotes");
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
() -> new PaneLocator(h).tabLabelsByTabId(),
() -> new PaneLocator(h).workspaceLabelsByWorkspaceId(), _ -> true,
callers::boundToArchitectSlot, callers.sendableObserverTarget());
String out = textOf(listFleetAsObserver(h, callers.leads(),
Map.of("term_collab_pane", "ops"), panes));
assertTrue(out.contains("\"panes\":["), out);
assertTrue(out.contains("\"sessionId\":\"term_sendable\""),
"an ordinary unclassified pane must still be sendable and visible: " + out);
assertFalse(out.contains("term_lead_pane"), "a lead's pane must not be enumerated: " + out);
assertFalse(out.contains("term_member_pane"), "a spawned member's pane must not be enumerated: " + out);
assertFalse(out.contains("term_collab_pane"), "a collaborator's pane must not be enumerated: " + out);
assertFalse(out.contains("term_architect_pane"),
"an unoccupied architect-slot pane must not be enumerated: " + out);
assertFalse(out.contains("\"paneId\""), "an observer's row must never carry paneId: " + out);
assertFalse(out.contains("\"workspaceId\""), "an observer's row must never carry workspaceId: " + out);
assertFalse(out.contains("\"workspaceLabel\""),
"an observer's row must never carry workspaceLabel: " + out);
assertFalse(out.contains("\"tabId\""), "an observer's row must never carry tabId: " + out);
assertFalse(out.contains("\"agentType\""), "an observer's row must never carry agentType: " + out);
assertFalse(out.contains("\"cwd\""), "an observer's row must never carry cwd: " + out);
}
/**
* Control for the test above: a primary's {@code panes} row is unchanged by fleetd #758 —
* {@code callerIsObserver} false keeps every field, including {@code paneId} and a spawned
* member's {@code cwd}.
*/
@Test
void listKeepsTheFullPaneRowForAPrimaryIncludingCwdAndPaneId() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
MemberSession spawned = sessions.acquire("ltms-local", "/worktree/member-1", null, null);
h.withAgent("claude-member", spawned.terminalId(), "w9:pMember", "w9:tMember");
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(Map::of, Map::of, _ -> true, _ -> false, _ -> false);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
Map.of(), "", Map.of(), false,
FleetMcp.CoordinationSource.none(), true, true, true, panes, true, false);
String out = textOf(res);
assertTrue(out.contains("\"sessionId\":\"" + spawned.terminalId() + "\""), out);
assertTrue(out.contains("\"paneId\":\"w9:pMember\""),
"a primary must still see the fleet_stop handle: " + out);
assertTrue(out.contains("\"cwd\":\"/worktree/member-1\""),
"a primary must still see a spawned member's worktree path: " + out);
assertTrue(out.contains("\"role\":\"dev\""), out);
}
/**
* fleetd #421: {@code mailbox.pending} counts only broker-ready messages, so a blocked lead's
* normal, healthy state is {@code "pending": 0} next to a non-empty {@code held[]} — which
@@ -2013,8 +2261,10 @@ class FleetMcpTest {
/**
* {@code fleet_status}'s pending-ask block (the question, its {@code turnId} and its ticket)
* 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.
* is shown only to the caller whose owner key created the delegation. An unnamed primary is
* held to the same rule: its owner key is {@code null}, which here does not match the named
* worker that created the delegation, so it sees none of the pending-ask fields either — the
* same as any other non-creating caller.
*/
@Test
void statusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
@@ -2054,8 +2304,13 @@ class FleetMcpTest {
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?"),
"a caller with no terminal (the unnamed primary) must see the question: " + unnamed);
assertTrue(unnamed.startsWith("idle"), "the base status must still be shown: " + unnamed);
assertFalse(unnamed.contains("which config file?"),
"an unnamed primary must not see a question on a delegation a named worker created: " + unnamed);
assertFalse(unnamed.contains(asking.turnId()),
"a non-creating unnamed primary must not see the turnId: " + unnamed);
assertFalse(unnamed.contains(ticket),
"a non-creating unnamed primary must not see the ticket: " + unnamed);
// Clean up the still-open ask so the background thread does not linger past the test.
String turnId = asking.turnId();
@@ -40,7 +40,7 @@ class LeadCoordLoopTest {
@Test
void deliversAHeldMessageToTheLeadPaneAndAcksIt() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
var herdr = new FakeHerdr().agentStatus("idle");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
@@ -57,7 +57,7 @@ class LeadCoordLoopTest {
@Test
void redeliveryOfAMessageAlreadyWrittenToThePaneIsAckedWithoutAnotherPaneWrite() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
var herdr = new FakeHerdr().agentStatus("idle");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
loop.tick();
@@ -71,7 +71,7 @@ class LeadCoordLoopTest {
@Test
void aRedeliveryIsAckedEvenWhileTheLeadIsMidTurn() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
var herdr = new FakeHerdr().agentStatus("idle");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
loop.tick();
@@ -91,7 +91,7 @@ class LeadCoordLoopTest {
@Test
void leavesTheMessageUnackedWhenTheLeadIsMidTurn() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
var herdr = new FakeHerdr().agentStatus("working");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("working");
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
@@ -103,7 +103,7 @@ class LeadCoordLoopTest {
@Test
void leavesTheMessageUnackedWhenNoLeadPaneIsKnown() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
var herdr = new FakeHerdr().agentStatus("idle");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
loop(channel, herdr, Map.of()).tick();
@@ -115,7 +115,8 @@ class LeadCoordLoopTest {
@Test
void leavesTheMessageUnackedWhenHerdrRefusesTheInjection() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
var herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("agent_not_found");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET)
.agentStatus("idle").agentSendFailsWith("agent_not_found");
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
@@ -126,7 +127,7 @@ class LeadCoordLoopTest {
@Test
void resolvesTheLeadByNameWhenSeveralAreKnown() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
var herdr = new FakeHerdr().agentStatus("idle");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
// Two leads on this daemon; only one carries the coord-id the mailbox is owned as.
var leads = new java.util.LinkedHashMap<String, String>();
leads.put("term_other", "some-other-lead");
@@ -142,7 +143,7 @@ class LeadCoordLoopTest {
@Test
void holdsWhenSeveralLeadsAreKnownAndNoneCarriesTheCoordId() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
var herdr = new FakeHerdr().agentStatus("idle");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
var leads = new java.util.LinkedHashMap<String, String>();
leads.put("term_one", "lead-one");
leads.put("term_two", "lead-two");
@@ -159,7 +160,7 @@ class LeadCoordLoopTest {
var channel = new FakeLeadChannel(SELF)
.hold(new LeadMessage("m1", PEER, SELF, "first"))
.hold(new LeadMessage("m2", PEER, SELF, "second"));
var herdr = new FakeHerdr().agentStatus("idle");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
loop.tick();
@@ -175,10 +176,51 @@ class LeadCoordLoopTest {
@Test
void anEmptyMailboxNeverTouchesHerdr() {
var channel = new FakeLeadChannel(SELF);
var herdr = new FakeHerdr().agentStatus("idle");
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
assertEquals(0, herdr.calls.size(), "an idle fleet must not poll a pane's status every tick");
}
@Test
void aLeadWithUnsubmittedTextInItsPromptBoxKeepsTheMessageHeldAndUnacked() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
var herdr = new FakeHerdr().detectionText(FakeHerdr.DRAFTED_PROMPT_CARET).agentStatus("idle");
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
assertEquals(0, prompts(herdr).size(),
"delivery pastes and submits, so it must not land on a half-typed line");
assertEquals(List.of(), channel.acked(), "an undelivered message stays on the broker");
assertFalse(channel.peek().isEmpty(), "and is still held");
}
@Test
void aMessageHeldForADraftIsDeliveredOnALaterTick() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
var herdr = new FakeHerdr().detectionText(FakeHerdr.DRAFTED_PROMPT_CARET).agentStatus("idle");
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
loop.tick();
assertEquals(0, prompts(herdr).size());
herdr.detectionText(FakeHerdr.IDLE_PROMPT_CARET);
loop.tick();
assertEquals(1, prompts(herdr).size(), "the held message lands once the box is empty");
assertEquals(List.of("m1"), channel.acked());
}
@Test
void anUnreadablePaneKeepsTheMessageHeld() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
var herdr = new FakeHerdr().detectionText("garbled ansi noise with no input box").agentStatus("idle");
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
assertEquals(0, prompts(herdr).size(), "a pane whose box cannot be found may be holding a draft");
assertEquals(List.of(), channel.acked());
}
}
@@ -7,6 +7,7 @@ import ch.qos.logback.core.read.ListAppender;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.lead.LeadContextGauge;
@@ -695,6 +696,8 @@ class LeadHeartbeatLoopTest {
private final List<String> sentTexts = new ArrayList<>();
private final List<String> promptTargets = new ArrayList<>();
private boolean throwOnNextSend = false;
/** What {@code agent.read} reports — the loop reads the lead's input box before it nudges. */
private String paneTail = FakeHerdr.IDLE_PROMPT_CARET;
FailableHerdrClient(String lead) {
this.lead = lead;
@@ -704,6 +707,10 @@ class LeadHeartbeatLoopTest {
throwOnNextSend = true;
}
void paneTail(String tail) {
this.paneTail = tail;
}
List<String> sentTexts() {
return List.copyOf(sentTexts);
}
@@ -722,6 +729,10 @@ class LeadHeartbeatLoopTest {
.put("terminal_id", lead)
.put("agent_status", "idle"));
}
if ("agent.read".equals(method)) {
return MAPPER.createObjectNode()
.set("read", MAPPER.createObjectNode().put("text", paneTail));
}
if ("agent.prompt".equals(method)) {
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
if (throwOnNextSend) {
@@ -738,4 +749,49 @@ class LeadHeartbeatLoopTest {
public void close() {
}
}
// ── the operator's own prompt box ────────────────────────────────────────────────────────────
@Test
void tickHoldsTheNudgeWhileTheLeadsPromptBoxHoldsUnsubmittedText() {
var herdr = new FailableHerdrClient(LEAD);
herdr.paneTail(FakeHerdr.DRAFTED_PROMPT_CARET);
var now = new AtomicLong(NOW);
@SuppressWarnings("unchecked")
List<MemberSession>[] rosterBox = new List[]{List.of()};
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
loop.tick(); // opens the idle window
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
loop.tick(); // would INJECT, but the operator is mid-sentence
assertEquals(0, herdr.sentTexts().size(),
"a nudge pastes and submits, so it must not land on a half-typed line");
herdr.paneTail(FakeHerdr.IDLE_PROMPT_CARET);
loop.tick();
assertEquals(1, herdr.sentTexts().size(), "the held nudge lands once the box is empty");
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"),
"and it still carries the notice the held tick did not spend: " + herdr.sentTexts().get(0));
}
@Test
void anUnreadablePaneHoldsTheHeartbeatNudge() {
var herdr = new FailableHerdrClient(LEAD);
herdr.paneTail("garbled ansi noise with no input box");
var now = new AtomicLong(NOW);
@SuppressWarnings("unchecked")
List<MemberSession>[] rosterBox = new List[]{List.of()};
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
loop.tick();
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
loop.tick();
assertEquals(0, herdr.sentTexts().size(),
"a pane whose box cannot be found may be holding a draft");
}
}
@@ -948,7 +948,7 @@ class MessageServiceTest {
}
@Test
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
void unnamedPrimaryIsRefusedFromANamedLeadsTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null,
Principal.leader("opus", "term_lead", 1));
awaitWaiting();
@@ -956,8 +956,51 @@ class MessageServiceTest {
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
assertNotNull(view, "the unnamed primary must be able to read any ticket");
MessageService.TaskView refused = messages.poll(ticket, Principal.primary(1).ownerKey());
assertNotNull(refused, "a different owner gets a refusal, not silence");
assertEquals(MessageService.Phase.FAILED, refused.phase());
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
assertNull(refused.reply(), "a refusal must never carry the reply text");
assertFalse(String.valueOf(refused).contains("primary-visible result"),
"the reply text must not appear anywhere in the refused view");
}
/**
* Positive control for {@link #unnamedPrimaryIsRefusedFromANamedLeadsTicket}: without this,
* that test would pass just as well if {@code poll} refused every caller.
*/
@Test
void unnamedPrimaryReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, Principal.primary(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");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, Principal.primary(2).ownerKey());
assertNotNull(view, "an unnamed primary must be able to read a ticket another unnamed primary created");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@Test
void theOneArgPollOverloadBypassesOwnershipEntirely() throws Exception {
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");
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 2000;
while (view == null || view.phase() != MessageService.Phase.DONE) {
if (System.currentTimeMillis() >= deadline) break;
view = messages.poll(ticket); // the one-arg, no-check overload -- no caller owner key at all
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(view, "the internal bypass must read a ticket owned by a named lead");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@@ -1978,7 +2021,9 @@ class MessageServiceTest {
/**
* A caller's owner key must match the key that created the delegation to see its pending
* question. The unnamed primary always sees it.
* question. The unnamed primary is held to the same rule as everyone else: its key is
* {@code null}, which here does not match the named worker that created this delegation, so
* it is refused too.
*/
@Test
void pendingAskGatesTheQuestionByTheDelegationsCreatorOwner() throws Exception {
@@ -1998,10 +2043,65 @@ class MessageServiceTest {
assertNotNull(own, "the creating caller must see its own open question");
assertEquals("which config file?", own.question());
MessageService.PendingAsk unnamed = messages.pendingAsk(T, null);
assertNotNull(unnamed, "a caller with no terminal (the unnamed primary) must always see the question");
assertNull(messages.pendingAsk(T, Principal.primary(1).ownerKey()),
"an unnamed primary must not see a question on a delegation a named worker created");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.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());
}
/**
* Positive control for {@link #pendingAskGatesTheQuestionByTheDelegationsCreatorOwner}:
* without this, that test's refusal would pass just as well if {@code pendingAsk} refused
* every caller. Here the delegation's creator is itself an unnamed primary (owner key
* {@code null}), so another unnamed primary's {@code null} key must still match it.
*/
@Test
void unnamedPrimarySeesItsOwnDelegationsPendingQuestion() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, Principal.primary(1));
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.PendingAsk unnamed = messages.pendingAsk(T, Principal.primary(2).ownerKey());
assertNotNull(unnamed, "an unnamed primary must see the question on a delegation another unnamed primary created");
assertEquals("which config file?", unnamed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(unnamed.turnId(), "config.yaml", 5000, null));
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());
}
/**
* {@code pendingAsk} has no public no-check overload the way {@link MessageService#poll}
* does, so this drives {@link MessageService#INTERNAL_NO_OWNER_CHECK} directly — the only way
* to exercise the bypass for this method.
*/
@Test
void pendingAskInternalBypassSeesAnyDelegationsPendingQuestion() throws Exception {
Principal creator = Principal.worker("term_creator", 1);
String ticket = messages.sendAsync(T, "task that asks", null, creator);
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.PendingAsk bypassed = messages.pendingAsk(T, MessageService.INTERNAL_NO_OWNER_CHECK);
assertNotNull(bypassed, "the internal bypass must see a question on a delegation a named worker created");
assertEquals("which config file?", bypassed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
@@ -2037,6 +2137,118 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
// --- outstanding(): fleet_handover{open}'s list of open tickets and asks -------------------
@Test
void outstandingReportsOwnedPendingTicketsWithPhases() {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket1 = messages.sendAsync(T, "task one", null, lead);
String ticket2 = messages.sendAsync(T, "task two", null, lead);
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
assertEquals(2, outstanding.tickets().size());
assertTrue(outstanding.tickets().stream().anyMatch(t ->
ticket1.equals(t.ticket()) && t.phase() == MessageService.Phase.PENDING && T.equals(t.target())),
"ticket1 must be reported PENDING: " + outstanding.tickets());
assertTrue(outstanding.tickets().stream().anyMatch(t ->
ticket2.equals(t.ticket()) && t.phase() == MessageService.Phase.PENDING && T.equals(t.target())),
"ticket2 must be reported PENDING: " + outstanding.tickets());
assertTrue(outstanding.asks().isEmpty(), "neither ticket has an open question");
}
@Test
void outstandingReportsAnOpenAsksTurnId() throws Exception {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "task that asks", null, lead);
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
assertEquals(1, outstanding.tickets().size());
assertEquals(MessageService.Phase.ASKING, outstanding.tickets().get(0).phase());
assertEquals(1, outstanding.asks().size());
MessageService.OutstandingAsk openAsk = outstanding.asks().get(0);
assertEquals(ticket, openAsk.ticket());
assertEquals(asking.turnId(), openAsk.turnId());
assertEquals(T, openAsk.workerSession());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, lead.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());
}
/**
* Positive control: {@code leadA} must still see its own ticket and ask, so {@code leadB}
* seeing neither is the owner filter at work and not {@code outstanding} refusing everyone.
*/
@Test
void outstandingDoesNotLeakAcrossNamedLeads() throws Exception {
Principal leadA = Principal.leader("opus", "term_a", 1);
Principal leadB = Principal.leader("sol", "term_b", 2);
String ticket = messages.sendAsync(T, "task that asks", null, leadA);
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Outstanding seenByA = messages.outstanding(leadA.ownerKey());
assertEquals(1, seenByA.tickets().size(), "lead A must see its own ticket");
assertEquals(1, seenByA.asks().size(), "lead A must see its own open ask");
MessageService.Outstanding seenByB = messages.outstanding(leadB.ownerKey());
assertTrue(seenByB.tickets().isEmpty(), "lead B must not see lead A's ticket");
assertTrue(seenByB.asks().isEmpty(), "lead B must not see lead A's open ask");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, leadA.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());
}
@Test
void outstandingIsEmptyCollectionsNotNullForACallerWithNothing() {
MessageService.Outstanding outstanding =
messages.outstanding(Principal.leader("opus", "term_lead", 1).ownerKey());
assertNotNull(outstanding.tickets(), "a caller with no delegations still gets a list, not null");
assertNotNull(outstanding.asks(), "a caller with no delegations still gets a list, not null");
assertTrue(outstanding.tickets().isEmpty());
assertTrue(outstanding.asks().isEmpty());
}
/**
* A ticket is destroyed on a timer once it goes terminal ({@link #pruneTerminalTickets}'s TTL
* runs from completion), so it is the one case where carrying the id forward actually matters
* — a still-PENDING ticket is in no such danger, its worker is still running.
*/
@Test
void outstandingReportsACompletedUncollectedTicketWithATerminalPhase() throws Exception {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "long task", null, lead);
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "async result"), "a reply resolves the async send");
awaitTicketPhase(ticket, MessageService.Phase.DONE);
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
assertEquals(1, outstanding.tickets().size());
MessageService.OutstandingTicket done = outstanding.tickets().get(0);
assertEquals(ticket, done.ticket());
assertEquals(MessageService.Phase.DONE, done.phase());
assertEquals(T, done.target());
}
// --- fleetd #715: answer() is gated on the caller that owns the turn -----------------------
/**
@@ -2164,7 +2376,7 @@ class MessageServiceTest {
private PushWiring wireWithPushLoop(int maxReminders, long backoffMs, java.util.function.LongSupplier nowNanos) {
PrimaryRegistry registry = new PrimaryRegistry(null);
registry.recordDelegation(T, LEAD);
FakeHerdr leadHerdr = new FakeHerdr();
FakeHerdr leadHerdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET);
AgentControl leadAgents = new AgentControl(leadHerdr);
var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor();
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
@@ -2201,7 +2413,7 @@ class MessageServiceTest {
java.util.function.LongSupplier nowNanos) {
PrimaryRegistry registry = new PrimaryRegistry(null);
registry.recordDelegation(T, LEAD);
FakeHerdr leadHerdr = new FakeHerdr();
FakeHerdr leadHerdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET);
AgentControl leadAgents = new AgentControl(leadHerdr);
ManualScheduler scheduler = new ManualScheduler();
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
@@ -2460,7 +2672,8 @@ class MessageServiceTest {
var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor();
// A backoff far longer than the test: the schedule is started but no tick ever fires, so
// decide() is read directly and nothing here depends on timing.
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(new FakeHerdr()), inbox,
ReplyPushLoop pushLoop = new ReplyPushLoop(registry,
new AgentControl(new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET)), inbox,
scheduler, 5, 60_000);
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop);
try {
@@ -3,6 +3,7 @@ package dev.ltms.fleet.msg;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.mcp.PrimaryRegistry;
@@ -1079,6 +1080,76 @@ class ReplyPushLoopTest {
"hitting the ticket reminder cap must count as exhausted");
}
// --- the operator's own prompt box ----------------------------------------------------------
@Test
void aLeadWithUnsubmittedTextInItsPromptBoxIsNotNudged() {
var herdr = new PaneTextHerdrClient(FakeHerdr.DRAFTED_PROMPT_CARET);
agents = new AgentControl(herdr);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(5, 100_000);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0),
"a nudge pastes and submits, so an idle lead mid-sentence must not be nudged");
assertEquals(0, herdr.promptCount(), "nothing reached the pane");
assertFalse(inbox.peek(WORKER).isEmpty(), "and the reply is still waiting to be collected");
}
@Test
void theSameLeadIsNudgedOnceItsPromptBoxIsEmpty() {
var herdr = new PaneTextHerdrClient(FakeHerdr.DRAFTED_PROMPT_CARET);
agents = new AgentControl(herdr);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(5, 100_000);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
herdr.paneText(FakeHerdr.IDLE_PROMPT_CARET);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"the box emptied, so the held nudge is due");
}
@Test
void anUnrecognisablePaneHoldsTheNudge() {
var herdr = new PaneTextHerdrClient("garbled ansi noise with no input box");
agents = new AgentControl(herdr);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(5, 100_000);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0),
"a pane whose box cannot be found may be holding a draft");
assertEquals(0, herdr.promptCount());
}
@Test
void aFailedPaneReadHoldsTheNudge() {
var herdr = new PaneTextHerdrClient(FakeHerdr.IDLE_PROMPT_CARET).failReads();
agents = new AgentControl(herdr);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(5, 100_000);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0),
"an unreadable box is treated as a draft, never as an empty one");
assertEquals(0, herdr.promptCount());
}
@Test
void aNudgeHeldForADraftIsSentOnALaterTick() throws Exception {
var herdr = new PaneTextHerdrClient(FakeHerdr.DRAFTED_PROMPT_CARET);
agents = new AgentControl(herdr);
loop(5, 50).onTicketTerminal("task-1", WORKER, false);
Thread.sleep(300);
assertEquals(0, herdr.promptCount(), "every tick holds while the operator is typing");
herdr.paneText(FakeHerdr.IDLE_PROMPT_CARET);
assertTrue(herdr.sendLatch.await(3, TimeUnit.SECONDS),
"the nudge lands on the first tick after the box empties");
}
// --- helpers -------------------------------------------------------------------------------
private ReplyPushLoop loop() {
@@ -1097,6 +1168,15 @@ class ReplyPushLoopTest {
return new AgentControl(new FakeHerdrClient(status));
}
/**
* The {@code agent.read} frame every fake here returns: a lead settled at an empty input box. The
* loop reads the box before it nudges, so a fake that answered nothing would read as a pane it
* cannot classify and hold every nudge.
*/
private static JsonNode emptyPromptBoxRead() {
return MAPPER.createObjectNode().set("read", MAPPER.createObjectNode().put("text", FakeHerdr.IDLE_PROMPT_CARET));
}
/** Non-recording (single-threaded) fake — safe for decide() tests. */
private static final class FakeHerdrClient implements HerdrClient {
private final String agentStatus;
@@ -1113,6 +1193,9 @@ class ReplyPushLoopTest {
.put("terminal_id", PRIMARY)
.put("agent_status", agentStatus));
}
if ("agent.read".equals(method)) {
return emptyPromptBoxRead();
}
return MAPPER.createObjectNode();
}
@@ -1142,6 +1225,9 @@ class ReplyPushLoopTest {
calls.add(Map.entry(method, params));
sendLatch.countDown();
}
if ("agent.read".equals(method)) {
return emptyPromptBoxRead();
}
return MAPPER.createObjectNode();
}
@@ -1179,6 +1265,9 @@ class ReplyPushLoopTest {
if ("agent.prompt".equals(method)) {
sendLatch.countDown();
}
if ("agent.read".equals(method)) {
return emptyPromptBoxRead();
}
return MAPPER.createObjectNode();
}
@@ -1216,6 +1305,9 @@ class ReplyPushLoopTest {
calls.add(Map.entry(method, params));
sendLatch.countDown();
}
if ("agent.read".equals(method)) {
return emptyPromptBoxRead();
}
return MAPPER.createObjectNode();
}
@@ -1266,6 +1358,9 @@ class ReplyPushLoopTest {
promptTargets.add(String.valueOf(p.get("target")));
sendLatch.countDown();
}
if ("agent.read".equals(method)) {
return emptyPromptBoxRead();
}
return MAPPER.createObjectNode();
}
@@ -1316,6 +1411,9 @@ class ReplyPushLoopTest {
if ("agent.prompt".equals(method)) {
promptTargets.add(String.valueOf(p.get("target")));
}
if ("agent.read".equals(method)) {
return emptyPromptBoxRead();
}
return MAPPER.createObjectNode();
}
@@ -1363,6 +1461,9 @@ class ReplyPushLoopTest {
promptTargets.add(String.valueOf(p.get("target")));
sendLatch.countDown();
}
if ("agent.read".equals(method)) {
return emptyPromptBoxRead();
}
return MAPPER.createObjectNode();
}
@@ -1374,4 +1475,60 @@ class ReplyPushLoopTest {
public void close() {
}
}
/**
* Thread-safe fake that always reports {@code idle} and serves a mutable pane tail, so a test can
* change what the lead's input box holds between ticks. Records every {@code agent.prompt}.
*/
private static final class PaneTextHerdrClient implements HerdrClient {
private final List<Map.Entry<String, Object>> prompts =
Collections.synchronizedList(new ArrayList<>());
private volatile String paneText;
private volatile boolean failReads = false;
volatile CountDownLatch sendLatch = new CountDownLatch(1);
PaneTextHerdrClient(String paneText) {
this.paneText = paneText;
}
void paneText(String text) {
this.paneText = text;
}
PaneTextHerdrClient failReads() {
this.failReads = true;
return this;
}
int promptCount() {
return prompts.size();
}
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
return MAPPER.createObjectNode()
.set("agent", MAPPER.createObjectNode()
.put("terminal_id", PRIMARY)
.put("agent_status", "idle"));
}
if ("agent.read".equals(method)) {
if (failReads) {
throw new HerdrException("herdr socket read timed out");
}
return MAPPER.createObjectNode()
.set("read", MAPPER.createObjectNode().put("text", paneText));
}
if ("agent.prompt".equals(method)) {
prompts.add(Map.entry(method, params));
sendLatch.countDown();
}
return MAPPER.createObjectNode();
}
@Override
public void close() {
}
}
}
@@ -91,8 +91,13 @@ class FleetAppAuthTest {
MessageService messages = new MessageService(agents, injector, new Rendezvous());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
// term_a is the only herdr-owned pane this fixture's PID can resolve to (FakeHerdr's canned
// pane list), and this helper's own contract above says that pane is the worker -- so it
// must be recognised as a live spawned member here, the same way a real roster would,
// rather than falling through to the observer floor.
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, tokenMode, token,
Map::of, new MemberRegistry(null));
Map::of, new MemberRegistry(null), t -> "term_a".equals(t) ? MemberRole.DEV : null,
Map::of);
metrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
app = new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
@@ -207,14 +212,15 @@ class FleetAppAuthTest {
}
/**
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket,
* while the creating worker and the unnamed primary both still read it. The ticket is minted
* directly on the shared {@link MessageService}, the same way {@code MessageServiceTest}
* drives {@link MessageService#poll(String, String)}, so this exercises only the REST poll
* route's own handling of the ownership already recorded on the ticket.
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket, and
* refuses an unnamed primary just the same: a named worker's ticket is not anyone else's to
* read, caller rank included. The ticket is minted directly on the shared
* {@link MessageService}, the same way {@code MessageServiceTest} drives
* {@link MessageService#poll(String, String)}, so this exercises only the REST poll route's
* own handling of the ownership already recorded on the ticket.
*/
@Test
void restPollRefusesADifferentWorkerButAllowsTheCreatorAndTheUnnamedPrimary() throws Exception {
void restPollRefusesADifferentWorkerAndAnUnnamedPrimaryOnANamedWorkersTicket() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
@@ -241,8 +247,10 @@ class FleetAppAuthTest {
HttpResponse<String> primary = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, primary.statusCode());
assertFalse(primary.body().contains("forbidden"),
"the unnamed primary must read any ticket: " + primary.body());
assertTrue(primary.body().contains("forbidden"),
"an unnamed primary must not read a ticket a named worker created: " + primary.body());
assertFalse(primary.body().contains("\"reply\""),
"a refusal must never carry reply text: " + primary.body());
} finally {
creatorApp.stop();
otherWorkerApp.stop();
@@ -250,6 +258,33 @@ class FleetAppAuthTest {
}
}
/**
* Positive control for
* {@link #restPollRefusesADifferentWorkerAndAnUnnamedPrimaryOnANamedWorkersTicket}: without
* this, that test's refusal would pass just as well if the route refused every caller. Here
* the ticket's creator is itself an unnamed primary, so another unnamed primary reading it
* over REST must still succeed.
*/
@Test
void restPollAllowsAnUnnamedPrimaryItsOwnTicket() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, new Rendezvous());
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
try {
String ticket = messages.sendAsync("term_a", "long task", null, Principal.primary(FakeHerdr.WORKER_PID));
HttpResponse<String> own = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, own.statusCode());
assertFalse(own.body().contains("forbidden"),
"an unnamed primary must read a ticket another unnamed primary created: " + own.body());
} finally {
primaryApp.stop();
}
}
/**
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
* caller's own terminal on the ticket it returns, so that caller can still poll its own
@@ -290,9 +325,10 @@ 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 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.
* {@code turnId} and its ticket only to the caller whose owner key created that delegation. An
* unnamed primary is held to the same rule: its owner key is {@code null}, which here does not
* match the named worker that created the delegation, so it sees none of the pending-ask
* fields either — the same as any other non-creating caller.
*/
@Test
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
@@ -321,7 +357,8 @@ class FleetAppAuthTest {
MessageService.TaskView asking;
deadline = System.currentTimeMillis() + 3000;
do {
asking = messages.poll(ticket, null);
// the no-check overload: this is test plumbing waiting for ASKING, not the gate under test
asking = messages.poll(ticket);
Thread.sleep(5);
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
@@ -342,8 +379,11 @@ class FleetAppAuthTest {
JsonNode primary = mapper.readTree(
send(primaryApp.port(), "GET", "/sessions/term_target/status", null, null).body());
assertEquals("which config file?", primary.get("question").asText(),
"a caller with no terminal (the unnamed primary) must see the question");
assertEquals("idle", primary.get("status").asText(), "the base status must still be shown");
assertFalse(primary.has("question"),
"an unnamed primary must not see a question on a delegation a named worker created: " + primary);
assertFalse(primary.has("turnId"), "a non-creating unnamed primary must not see the turnId: " + primary);
assertFalse(primary.has("ticket"), "a non-creating unnamed primary must not see the ticket: " + primary);
// Clean up the still-open ask so the background thread does not linger past the test.
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
@@ -398,7 +438,8 @@ class FleetAppAuthTest {
MessageService.TaskView asking;
deadline = System.currentTimeMillis() + 3000;
do {
asking = messages.poll(ticket, null);
// the no-check overload: this is test plumbing waiting for ASKING, not the gate under test
asking = messages.poll(ticket);
Thread.sleep(5);
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
@@ -223,6 +223,37 @@ class FleetAppTest {
assertTrue(body.has("detail"), res.body());
}
@Test
void agentsReportsTheAgentsTabLabel() throws Exception {
FakeHerdr herdr = new FakeHerdr().withTab("w2", "w2:t7", "trinotes");
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "GET", "/agents");
assertEquals(200, res.statusCode(), res.body());
JsonNode agents = mapper.readTree(res.body()).get("agents");
assertEquals(1, agents.size());
assertEquals("sess-1111", agents.get(0).get("sessionId").asText());
assertEquals("trinotes", agents.get(0).get("label").asText());
}
/**
* The tab-label scan ({@code workspace.list}/{@code tab.list}) is decoration on top of
* {@code workers.list()}'s own agent roster, so its failure must not cost that roster: a row
* reports a {@code null} label instead, never the {@code herdr_error} envelope.
*/
@Test
void agentsStillReportsTheRosterWhenTheLabelScanFails() throws Exception {
FakeHerdr herdr = new FakeHerdr().workspaceListFailsWith("unavailable");
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "GET", "/agents");
assertEquals(200, res.statusCode(), res.body());
JsonNode agents = mapper.readTree(res.body()).get("agents");
assertEquals(1, agents.size());
assertEquals("sess-1111", agents.get(0).get("sessionId").asText());
assertTrue(agents.get(0).get("label").isNull(), "a failed label scan must report a null label, not fail the roster: " + res.body());
}
@Test
void spawnWorkerLandsInOwnTabInWorkerSpaceAndInjectsBaseUrl() throws Exception {
FakeHerdr herdr = new FakeHerdr();
+3 -2
View File
@@ -1507,6 +1507,7 @@ else
grep -E ' (ERROR|SEVERE) ' "$FRESH_LOG" | tail -5 | sed 's/^/ /'
fi
echo
echo " Next: call fleet_whoami and confirm it still answers 'primary'. A lead whose tab label"
echo " no longer matches fleet.leaders.*.tab is demoted to worker and refuses orchestration."
echo " Next: call fleet_whoami and confirm it still answers 'primary'. A lead is found by its"
echo " tab being labelled 'lead' AND sitting in fleet.leaders.<name>.workspace; if either stops"
echo " matching, the lead is demoted to worker and refuses orchestration."
echo