Compare commits

..

16 Commits

Author SHA1 Message Date
Dai Ha 6022aef612 fleetd #778: let an observer read its own ticket, never nudge a role to run a call it is refused
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 59s
CI / build (pull_request) Failing after 2m15s
Part 1: Authz.TASK_READ now grants an observer the ticket it created (a new
observerOwnsTicket classifier), via MessageService.ownsTicket(String, String)
(side-effect-free) and a new Authz.permits overload. fleet_status and
fleet_poll{target} (DRAIN) stay refused for an observer. REST and MCP share
the same gate.

Part 2: ReplyPushLoop no longer nudges a reply/question delegator to run
fleet_poll(target=...) or fleet_send(turnId=...) unless Authz would actually
grant that call to its role. FleetMcp.sendHandler asks Authz.permits(DRAIN)/
permits(ANSWER) once, eagerly, while the real Principal is in scope, and
threads the booleans through PrimaryRegistry.recordDelegation into the new
Delegation fields and nudgeMayDrainFor/nudgeMayAnswerFor accessors --
plain booleans, not a Role reference, so no new or widened edge crosses the
frozen auth<->mcp<->msg package-cycle baseline. ReplyPushLoop filters a
reply/question out of its pending set when the grant is false, falling back
to the durable inbox / the ask window lapsing, same as when no nudge target
is known at all. Log lines that called a non-lead nudge target "lead" now
say "nudge target".
2026-10-06 13:55:26 +02:00
Dai Ha 66ab9cab24 fleetd #791: addendum — plugin 0.3.0 has no mount and is installed in four instances
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m59s
2026-10-06 08:03:50 +02:00
Dai Ha f611488c9b fleetd #791: README — the mod still reads FLEETD_MCP_URL as an optional override
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 56s
CI / build (push) Failing after 2m4s
2026-10-06 08:03:06 +02:00
Dai Ha 44d16e7639 fleetd #791: merge PR #794 — fleet plugin 0.3.0, no MCP mount, mod off for members
The plugin no longer carries .mcp.json; the instance or the project mounts
fleet. The mod asks fleet_whoami once and skips the fleet_inbox poll for a
worker or an architect, so members keep paste delivery. The mod reads
FLEETD_MCP_URL with $.env.get (documented in the mods API), defaulting to
http://127.0.0.1:8765/mcp.

Lead verification on 87f7c03: claude plugin validate plugin passed;
claude plugin test plugin 21 pass, 0 fail. With the role gate removed, the
worker and architect tests fail (19 pass, 2 fail), so they are not vacuous.
2026-10-06 08:03:01 +02:00
Dai Ha 87f7c03535 fleetd #791: fleet plugin 0.3.0 — drop the MCP mount, gate the mod off for members
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 2m7s
- Delete plugin/.mcp.json; mounting fleet is now the instance's or the
  project's job, not the plugin's. The setup skill still writes the
  project .mcp.json entry.
- register.js: read FLEETD_MCP_URL via $.env.get, defaulting to
  http://127.0.0.1:8765/mcp, documented at
  https://code.claude.com/docs/en/plugins/mods/api.md ("$.env — get
  and set environment variables").
- Gate the fleet_inbox poll by role: a worker or architect learns its
  role once from fleet_whoami and then never polls fleet_inbox; every
  other role polls as before. An unreachable daemon retries whoami on
  a later tick rather than deciding.
- Bump 0.2.0 -> 0.3.0 in plugin.json and marketplace.json, and update
  README/SKILL.md to drop every claim that the plugin mounts fleetd.
- Add 6 tests for the role gating; all 21 tests and
  `claude plugin validate` pass.
2026-10-06 06:46:26 +02:00
Dai Ha 562354a58d fleetd #790: document observer -> lead send in the bridge block
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 1m59s
Invariant 3, the observer paragraph and the unconfigured-pane row now say an
observer may send to a lead. Invariant 4 now says a direct send to a lead pane
is not box-gated yet (#793). Wiki template synced (check printed in sync: True).
2026-10-06 06:25:31 +02:00
Dai Ha eaa1f0d170 fleetd #790: merge PR #792 — an observer pane may fleet_send to a lead
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 56s
CI / build (push) Failing after 2m5s
An observer's SEND classifier (CallerResolver.observerSendTarget, renamed from
sendableObserverTarget) now accepts a configured lead terminal as well as
another observer pane. Collaborators, architect slots and spawned members stay
refused. The [fleet_send from observer term_…] prefix is kept. An observer
learns a lead's sessionId from its filtered fleet_list panes rows (role "lead").

Lead verification: mvn clean install on 72ea7de in a throwaway worktree,
exit 0, surefire totals 2208 run / 0 failures / 0 errors / 0 skipped.

Follow-up found in review: #793 (the Injector has no prompt-box gate for a
direct send to a lead pane).
2026-10-06 06:24:33 +02:00
Dai Ha 72ea7def0a fleetd #790: let an observer pane fleet_send to a lead
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 2m7s
An observer could only reach another observer pane, so a hand-opened tab
could answer a lead with fleet_reply but never start a conversation with
one.

CallerResolver.sendableObserverTarget() is renamed observerSendTarget()
and now accepts a configured lead terminal as well as a pane that falls
to the observer floor. It reads the same maps resolve() reads, in the
same order: a live spawned member is refused first, a lead is then
accepted even when the same pane is also bound to an architect slot,
and a collaborator or an architect-slot pane is refused. A collaborator,
an architect and a spawned member stay unreachable.

Authz keeps its one SEND case; only the classifier it is handed widened,
so the MCP gate (FleetMcp.denyFor) and the REST gate (FleetApp.allow)
give the same answer from the same predicate.

fleet_list shows the lead's address in the observer's filtered `panes`
rows rather than by adding the `leads` array: the panes filter already
reads the same classifier as the SEND gate, so reachability and
visibility cannot drift, and the reduced row carries no lead name,
context window or config dir.

Gates traced end to end for an observer -> lead send: denyFor, the
sendTool argument validation (profileTargetError rejects profile names
only), MessageService (records the caller's owner key, no role gate),
the injector's readiness gate (Fleetd.deliverableTo accepts a lead
through the leads map, never a presence entry) and its status gate.
Unchanged: an observer still holds no TASK_READ, so it cannot poll the
ticket a wait:false send returns.

Tests: observer -> lead allowed and observer -> collaborator / architect
slot / spawned member refused over denyFor, each with a positive control
in the same test; the same matrix over the REST route; a real MCP client
through the real MessageService and Injector to the real herdr call,
proving the [fleet_send from observer term_...] prefix reaches a lead's
pane and that the lead passes the production readiness gate; and the
observer's fleet_list panes rows now carrying the lead.

mvn clean install in fleetd/: Tests run: 2208, Failures: 0, Errors: 0,
Skipped: 0 — BUILD SUCCESS. Reverting the widening in CallerResolver
alone turns 7 of the new assertions red across all five touched test
classes, so none of them passes vacuously.
2026-10-06 06:20:43 +02:00
Dai Ha 98f3cd5aa7 fleetd #788: document fleet_inbox in the bridge block and the addendum
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 1m57s
Block: observer INBOX right, invariant 3 (inbox is only-as-itself),
invariant 4 (a pane that collects its own mail skips the paste and the
box gate, not the status gate), and an intent-table row. The wiki
template carries the same block; the sync check printed in sync: True.
2026-10-06 05:54:13 +02:00
Dai Ha aec5ff2c2f fleetd #788: merge PR #789 — fleet_inbox, so a pane can collect its own mail
A pane that calls fleet_inbox within 15 s is mod-served: the Injector
offers its next message for collection instead of pasting it, and keeps
it queued until the pane takes it. A pane that stops polling gets its
mail pasted again. The fleet mod polls fleet_inbox every 3 s on one
reused MCP session and hands each message to Claude with
$.prompt.submit.

Lead proof at d9fedfb, in a throwaway worktree:
mvn clean install exit 0; surefire XML 182 files, 2202 tests,
0 failures, 0 errors, 0 skipped. claude plugin validate plugin: passed.
claude plugin test plugin: 15 pass, 0 fail.
2026-10-06 05:53:16 +02:00
Dai Ha d9fedfbc5a fleetd #788: review fixes — cancel vs a collected offer, one MCP session
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 1m59s
Four items from the lead review of PR #789.

1. Injector.cancel touched the inbox offer unconditionally. A pane takes an
   offer on the MCP thread, so between its drain and the next poll the head
   entry is taken but still QUEUED. Cancelling it returned CANCELLED for text
   the pane already held, and cancelling a LATER entry nulled t.inboxOffer and
   so lost the only record that the head was taken, which made the next poll
   offer it a second time. cancel now touches the offer only when the entry is
   the offered head, withdraws first, and answers DELIVERED when the pane took
   it.

2. The mod sent initialize on every fleetTool call and never a DELETE, so
   fleetd's transport kept one session per call -- 20 a minute per pane at a
   3s poll. The mod now holds one MCP session and reopens it only on the 404
   "Session not found" fleetd answers for an id it no longer holds (measured
   against the live daemon, not assumed).

3. MessageService.collectInbox had been inserted between reply's javadoc and
   reply, leaving that block attached to nothing. Moved above it.

4. The timer's catch comment said a throw would kill the timer. It does not;
   the catch keeps every tick from writing an error to the debug log.

mvn clean install: Tests run: 2202, Failures: 0, Errors: 0, Skipped: 0,
BUILD SUCCESS. Counted again over 182 target/surefire-reports/TEST-*.xml:
2202/0/0/0. claude plugin validate plugin and claude plugin test plugin both
exit 0, 15 pass 0 fail.

Three mutation checks, each reverted:
- the exact pre-review cancel body: the two new cancel tests fail with
  "expected: <DELIVERED> but was: <CANCELLED>" and "expected: <true> but was:
  <false>", the other six pass;
- initialize on every call: the two session tests fail (1 vs 2 initializes);
- no 404 retry: only the retry test fails (2 vs 1 initializes).
2026-10-06 05:48:34 +02:00
Dai Ha ac535503ff fleetd #788: fleet_inbox, so a pane can collect its own mail
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Failing after 2m15s
A Claude Code session running the fleet mod now pulls its messages from
fleetd and submits them with $.prompt.submit, instead of having them typed
into its terminal by herdr. fleetd names a caller by its pane, so this is
the hop that works between two Claude accounts on one host.

New MCP tool fleet_inbox. It takes no arguments: the pane is the caller's
connection-resolved terminal, so no caller can read another pane's mail.
Authz gets an INBOX action, grouped with REPLY and ASK as only-as-itself.

A pane becomes mod-served by calling fleet_inbox, for a 15s window that
each call renews. The Injector asks that question at the moment it is
about to deliver, and offers the message for collection rather than typing
it. The message stays at the head of the Injector's queue until the pane
takes it, so:

  - delivery is recorded when the pane really has the text, not when it
    was offered, and fleet_poll reports the same states as for a typed
    message;
  - a pane that stops polling leaves the window and its mail is typed
    instead, with nothing stranded and nothing delivered twice. The offer
    is withdrawn under the same monitor that drains it, so an entry is
    removed exactly once;
  - a collected message gets no Enter nudge. Nothing was typed, and an
    Enter in a lead's pane would submit whatever its operator was writing.

cancel, drop and both grace expiries withdraw the offer too, so a message
the caller was told never arrived can never arrive later.

Mod side: a 3s timer collects through the existing fleetTool helper and
submits each message. A fleetd that is down returns from the timer rather
than throwing out of it, which would stop every later poll.

Build: mvn clean install in fleetd/ — Tests run: 2200, Failures: 0,
Errors: 0, Skipped: 0; BUILD SUCCESS. claude plugin validate plugin and
claude plugin test plugin both pass (13 pass, 0 fail).

Two mutations confirm the new tests bind: forcing isModServed false fails
9 of 15, and removing the stop-polling fallback fails exactly
aPaneThatStopsCollectingHasItsMailTypedInstead.
2026-10-06 05:31:08 +02:00
Dai Ha 654b3b5e14 fleetd #788: fleet mod with presence, mail and a fleetd identity probe
Turns the fleet plugin into a Claude Code mod (hooks/hooks.json +
register.js). /fleet-peers and /fleet-mail carry messages through
$.store; /fleet-whoami calls fleet_whoami on the local fleetd over
$.http.fetch.

Measured 2026-10-06 on Claude Code 2.1.290: $.store lives at
<CLAUDE_CONFIG_DIR>/plugins/store/, so the store mailbox does not cross
the ltms and work accounts (two inodes, different rows). /fleet-whoami
reached fleetd from both config dirs, so fleetd is the cross-account hop.
The fleetd pull path that finishes the adapter is #788.

claude plugin validate plugin: passed. claude plugin test plugin:
10 pass, 0 fail.
2026-10-06 05:09:04 +02:00
Dai Ha 1fae8b81a0 fleetd #782: tell a placeholder hint apart from real typed text in the prompt box
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 2m9s
The box gate held every delivery into an idle lead pane. Claude Code
draws a placeholder hint in the empty input box - the pane's own last
submitted prompt - and draws it faint. PromptBox read the pane with
source 'detection', which carries no escape codes at all, so the
dimness was gone before classify ran and the hint read as a draft.

- AgentControl gains readWithStyling, which sends strip_ansi: false.
  The two-argument read is untouched for its other callers.
- PROBE_SOURCE moves to 'visible', the source that does carry styling.
- markerLength skips leading SGR codes, because the grey drawn on an
  empty box's caret sits before the marker glyph.
- boxContent decodes the line into (char, faint) pairs and drops every
  character inside a faint span.
- boxContent also uses Character.isSpaceChar, not only isWhitespace.
  This was cosmetic on 'detection', which trimmed a trailing U+00A0;
  on 'visible' an empty box arrives as '<caret>\u00a0' and would
  otherwise read as DRAFT with 1 character.

Verified in a throwaway worktree at 4ef8156: mvn clean install, BUILD
SUCCESS, 179 test classes, 2184 tests, 0 failures, 0 errors, 0 skipped,
counted from fleetd/target/surefire-reports/*.xml. PromptBoxTest is 21
tests, up from 15.

Measurements behind the four fixtures, taken from seven live panes with
'herdr agent read <paneId> --source visible --ansi': the two held panes
carried SGR 2 around their hint text, the five empty boxes carried no
faint span, and the caret grey on one empty box was a 24-bit colour
sitting before the marker.

Two checks on the move to 'visible' that the fix depends on:

- The active-turn marker 'esc to interrupt' still arrives as one
  contiguous string with styling kept, so the UNREADABLE guard for a
  live turn is intact.
- Across all seven panes the only intensity codes emitted are 2 and 0.
  SGR 22 (normal intensity) never appears, and no compound code mixes
  an intensity with a colour, so exact-matching those two codes is
  enough for this build. A build that emitted 22 would leave faint on
  and could read a drafted box as empty; tracked on #782.

PR #784's HOLD_GIVE_UP_STREAK is deliberately not here. I asked for it
and it was wrong: the hold cadence measures 15.07s, so 1200 holds is
about five hours, and giving up submits whatever sits in the box.
2026-10-06 04:28:02 +02:00
Dai Ha 4ef815682b fleetd #782: tell a placeholder hint apart from real typed text in the prompt box
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 2m9s
PromptBox probed the detection region, which carries no ANSI styling, so a dim
placeholder hint (the pane's last submitted prompt) read as the operator's draft
and held the delivery forever.

Move the probe to the visible region, read with strip_ansi:false
(AgentControl.readWithStyling), skip leading SGR escapes before matching the
box marker, use Character.isSpaceChar so a non-breaking space in an empty box
counts as padding, and exclude any character drawn inside a faint (SGR 2) span
from the box content. An all-faint box now classifies as EMPTY instead of
UNREADABLE.
2026-10-06 04:22:59 +02:00
Dai Ha 6596458ce6 fleetd #780: queued sends no longer lie to fleet_poll or fleet_send
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 2m2s
Three false receipts on the delivery path to a pane are fixed:

- fleet_poll reported "pending — worker working" for a message still
  sitting in the injector queue. MessageService.pendingDetail now reads
  Injector.queuedWaitMillis and says "queued, not yet delivered".
- fleet_send's accept text carried no warning when the target was not
  injectable. FleetMcp.sendAsync now appends one.
- a queued message had no timeout at all. Injector gains
  QUEUE_WAIT_GRACE_POLLS (4800, ~20 min at the 250ms poll), mirroring
  READINESS_GRACE_POLLS, and fails the queue through onTurnFailed.

Verified in a throwaway worktree at 1cc0322: mvn clean install, BUILD
SUCCESS, 179 test classes, 2179 tests, 0 failures, 0 errors, 0 skipped,
counted from fleetd/target/surefire-reports/*.xml. All four new tests
are present in the XML and passed.

Reviewed by two reviewers. Two findings, both confirmed in the code and
both left open rather than fixed here:

- FleetMcp.notInjectableWarning reads only AgentStatus.injectable(), not
  the Fleetd.deliverableTo presence gate the Injector actually uses, so
  an idle pane that has not connected its bridge MCP still gets a clean
  accept receipt. This is the headline false receipt of #780 and is the
  reason #780 stays open.
- the queueStalled path does not call forget.accept. The proposed fix is
  rejected: a target stuck UNKNOWN may be alive but unclassifiable, and
  forgetting its presence would make it permanently undeliverable, since
  presence is only re-established by an MCP request an idle pane has no
  reason to make. Presence is already cleared on session release
  (SessionManager.java:477).
2026-10-06 04:03:43 +02:00
36 changed files with 2729 additions and 299 deletions
+2 -2
View File
@@ -8,8 +8,8 @@
{
"name": "fleet",
"source": "./plugin",
"description": "Mount the fleetd MCP gateway and apply standard Claude Code settings so a session can orchestrate delegated workers. Ships no credentials.",
"version": "0.2.0",
"description": "Apply standard Claude Code settings so a session can orchestrate delegated workers, and run the fleet mod for cross-session messaging. Mounting the fleetd MCP gateway is the instance's or the project's job, not this plugin's. Ships no credentials.",
"version": "0.3.0",
"author": {
"name": "LTMS"
}
+32 -11
View File
@@ -33,10 +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`, 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`.
to `READ`/`METRICS`, to `REPLY`/`ASK`/`INBOX` on its own pane, and to `SEND` only to a lead or to a target that
resolves as an observer too — never to a collaborator, an architect, 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`;
a lead's row reads `role: "lead"`.
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
@@ -67,9 +68,10 @@ and the sender silently receives nothing. Fail toward the recoverable error.
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,
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.
never to a spawned member's terminal, while an observer may send only to a lead or another observer
pane;
reply/ask/inbox are only-as-itself — any peer may answer, or collect its mail, 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
@@ -77,10 +79,19 @@ and the sender silently receives nothing. Fail toward the recoverable error.
**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
is clear, and a pane the daemon cannot read as a box waits too. A direct `fleet_send` to a lead's
pane does **not** wait for the box yet (fleetd #793); a lead that runs the fleet mod collects
that mail instead of having it pasted, so it is not exposed. 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.
**A pane that collects its own mail skips the paste.** A session running the fleet mod calls
`fleet_inbox` on a timer. While its last call is under 15s old, the daemon queues that pane's
next message for collection instead of pasting it, and the mod hands it to Claude with
`$.prompt.submit` — so the box gate does not apply, but the status gate still does. When the calls
stop, the pane's mail is pasted again, and nothing is lost. Pasted or collected, the fleet
channel also crosses Claude accounts on one host, because the daemon names a caller by its pane;
`SendMessage`, `ListAgents` and the mod's own store each stay inside one account.
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.
@@ -189,10 +200,11 @@ 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 |
| 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 you or to another observer pane, but never to a collaborator or a member. 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}` |
| Collect the mail queued for your OWN pane, instead of having it pasted | `fleet_inbox` — no arguments, any role, own pane only. The fleet mod calls it on a timer; you rarely call it by hand |
| Tear down a member | `fleet_stop{paneId}` |
| Replace your OWN lead session when its context is full | `fleet_handover{action:"open", reason?}` → write the handover file it names → `fleet_handover{action:"confirm", token, operatorConfirmed}`. Primary-only. **In that order**: the file must be modified *after* `open`, or `confirm` refuses it as stale. There is no terminal parameter — the pane is always your own, so you can never roll another lead. `{action:"cancel", token}` drops a pending request |
@@ -345,8 +357,17 @@ must obey belongs in the charter, not here.
`handover` (write the file a fresh lead session inherits when the outgoing one hands off,
fleetd #480).
- **This repo is also a Claude Code marketplace, and ships a plugin.** `.claude-plugin/marketplace.json`
points at `plugin/`, which carries the MCP mount and the `setup` skill
(`/claude-bridge:setup` — make any project bridge-ready). It was added in CB-527 and then went
points at `plugin/`, which carries the `setup` skill (`/fleet:setup` — make any project
bridge-ready) and no MCP mount: the instance or the project mounts `fleet`. `plugin/` is also a Claude Code mod (`plugin/hooks/`):
it polls `fleet_inbox` and hands each message to Claude with `$.prompt.submit`. Measured
2026-10-06: `fleet@fleetd` 0.3.0 is installed at user scope in the `gx10`, `ltms`, `ollama` and
`work` instances, so every session started from them runs the mod, and a spawned member on those
config dirs loads it too — but the mod skips the inbox poll for a `worker` or an `architect`, so a
member still gets its brief pasted. The install made a copy under each instance's
`plugins/cache/fleetd/fleet/0.3.0`, so assume an edit to `plugin/` reaches sessions only after a
version bump and `claude plugin update fleet@fleetd` per instance. Re-measure with
`grep -l '"fleet@fleetd"' ~/.ccs/instances/*/plugins/installed_plugins.json`; delete this sentence
if the plugin is uninstalled. It was added in CB-527 and then went
unmentioned by every instruction file, so it drifted and a later session planned it from scratch
(#362). **Read `plugin/` before designing anything about onboarding a project.** Two limits are
structural, not bugs: a plugin cannot carry the role agent files, because
@@ -32,6 +32,13 @@ public final class Authz {
REPLY,
/** A worker's mid-turn question to the primary. */
ASK,
/**
* Collect the messages queued for the caller's OWN pane, instead of having them typed into
* its terminal. Grouped with {@link #REPLY} and {@link #ASK} below as an only-as-itself
* action: the pane is always the caller's connection-resolved terminal, never an argument,
* so no caller can collect another pane's mail.
*/
INBOX,
/** Collect held replies from a session's inbox. */
DRAIN,
/** Read-only roster, profile, and identity observation: no ticket, task, or turn state. */
@@ -73,55 +80,88 @@ public final class Authz {
/**
* 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}.
* CallerResolver#observerSendTarget()} is the real one, read from the same maps {@code
* CallerResolver#resolve} consults, so a target that classifier accepts is one {@code resolve}
* would actually resolve as a lead ({@link Role#PRIMARY}) or as {@link Role#OBSERVER}.
*/
public static final Predicate<String> NO_KNOWN_OBSERVER_TARGET = target -> false;
public static final Predicate<String> NO_OBSERVER_SEND_TARGET = target -> false;
/**
* The fail-closed classifier for an observer's {@code TASK_READ}: answers no for every
* ticket, so the grant is refused unless a caller supplies a real one. {@code
* MessageService#ownsTicket(String, String)} is the real one, so a ticket this classifier
* accepts is one that caller's own {@code fleet_send(wait:false)} actually created.
*/
public static final Predicate<String> NO_OBSERVER_OWNED_TICKET = ticket -> false;
/**
* Convenience form for a caller with no classifier to supply. Fails closed: a collaborator's
* 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.
* or an observer's {@code SEND}, and an observer's {@code TASK_READ}, are refused, as if no
* terminal were a configured lead, collaborator, or observer-reachable target, and no ticket
* were the caller's own — the same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR},
* {@link #NO_OBSERVER_SEND_TARGET} and {@link #NO_OBSERVER_OWNED_TICKET} give explicitly.
* Every other action's result is identical to the full form's, since none of them consult
* any of the three classifiers.
*
* <p>Its default classifiers deny every collaborator and every observer, so a caller
* enforcing authorization must use the five-argument form instead.
* enforcing authorization must use the full form instead.
*/
public static boolean permits(Principal caller, Action action, String targetSession) {
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR, NO_KNOWN_OBSERVER_TARGET);
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR,
NO_OBSERVER_SEND_TARGET, NO_OBSERVER_OWNED_TICKET);
}
/**
* 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.
* {@code SEND}. An observer's {@code SEND} and {@code TASK_READ} still fail closed — a
* caller enforcing all three grants must use the full form.
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator) {
return permits(caller, action, targetSession, knownLeadOrCollaborator, NO_KNOWN_OBSERVER_TARGET);
return permits(caller, action, targetSession, knownLeadOrCollaborator,
NO_OBSERVER_SEND_TARGET, NO_OBSERVER_OWNED_TICKET);
}
/**
* As {@link #permits(Principal, Action, String, Predicate)}, with a real classifier for an
* observer's {@code SEND} too. An observer's {@code TASK_READ} still fails closed — a caller
* enforcing all three grants must use the full form.
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> observerSendTarget) {
return permits(caller, action, targetSession, knownLeadOrCollaborator, observerSendTarget,
NO_OBSERVER_OWNED_TICKET);
}
/**
* 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}), for a
* collaborator's {@code SEND}, and for an observer's
* {@code SEND}, ignored otherwise, may be {@code null}
* worker-scoped actions ({@code REPLY}, {@code ASK},
* {@code INBOX}), for a collaborator's {@code SEND}, for an
* observer's {@code SEND}, and — carrying a ticket id instead
* of a session id — for an observer's {@code TASK_READ},
* 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
* @param observerSendTarget whether a terminal is one this daemon would itself resolve as
* a lead ({@link Role#PRIMARY}) or as {@link Role#OBSERVER} —
* consulted only for an observer's {@code SEND}, to confine it
* to a lead or another observer pane and never a collaborator,
* an architect, or a spawned member
* @param observerOwnsTicket whether {@code targetSession} (here, a ticket id) was created
* by this same caller — consulted only for an observer's
* {@code TASK_READ}, to confine it to a ticket its own
* {@code fleet_send(wait:false)} created, never another
* caller's
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> knownObserverTarget) {
Predicate<String> observerSendTarget,
Predicate<String> observerOwnsTicket) {
if (caller == null || caller.isAnonymous()) {
return false; // authenticated as nothing ⇒ authorized for nothing
}
@@ -135,12 +175,12 @@ public final class Authz {
// 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.
// a target that would itself resolve as a lead or as another observer, never a
// collaborator, an architect, 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.isObserver() && knownObserverTarget.test(targetSession));
|| (caller.isObserver() && observerSendTarget.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:
@@ -159,7 +199,7 @@ public final class Authz {
// collaborator's own pane passes through the same check, so each can answer a funnel
// that delegated to it. An unnamed primary (token/loopback, no pane) owns nothing and
// is still excluded.
case REPLY, ASK -> caller.ownsSession(targetSession);
case REPLY, ASK, INBOX -> caller.ownsSession(targetSession);
// READ is roster, profile, and identity observation — fleet_list, fleet_profiles, and
// fleet_whoami — and carries no secrets: no ticket reply, no pending question, and no
@@ -170,11 +210,15 @@ public final class Authz {
|| caller.isCollaborator() || caller.isObserver();
// Ticket polling and session status, open to every role READ is open to except a
// collaborator or an observer. MessageService compares a ticket's creator to the
// caller on every read as well, so dropping this gate would not expose another
// session's reply — it would move the refusal later and widen what a caller that
// never orchestrates can probe.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// collaborator or an observer — with one exception: an observer may poll a ticket its
// own fleet_send(wait:false) created, confined by observerOwnsTicket. Session status
// is not ticket-scoped, so that call site supplies no real classifier here and an
// observer's TASK_READ on it stays refused. MessageService compares a ticket's
// creator to the caller on every read as well, so dropping this gate would not expose
// another session's reply — it would move the refusal later and widen what a caller
// that never orchestrates can probe.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect()
|| (caller.isObserver() && observerOwnsTicket.test(targetSession));
// fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect
// holds READ today (CB-548), so "not primary" must mean not-architect here too — this
@@ -271,22 +271,27 @@ public final class CallerResolver {
}
/**
* 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.
* Whether {@code target} names a terminal an observer may {@code SEND} to: one this resolver
* would itself resolve as a lead ({@link Role#PRIMARY}) or as {@link Role#OBSERVER}. Read from
* the same maps and functions {@link #resolve} consults, and in the same order, so a target
* this accepts is exactly one {@code resolve} would hand back one of those two roles for, and
* the reverse.
*
* <p>A live spawned member is refused first, whatever a tab map says about its terminal — the
* order {@link #resolve} itself uses. A pane named as a lead is then accepted even when it is
* also bound to an architect slot, because that is the role {@code resolve} gives it.
*/
public Predicate<String> sendableObserverTarget() {
public Predicate<String> observerSendTarget() {
return target -> target != null
&& spawnedMemberRole.apply(target) == null
&& !leadTerminals.get().containsKey(target)
&& !boundToArchitectSlot(target)
&& !collaboratorTerminals.get().containsKey(target);
&& (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
* architect — the one classifier {@link #observerSendTarget()} and {@code FleetMcp}'s
* {@code panes} row both read, so a pane's reported role and its {@code SEND} reachability can
* never drift apart.
*/
@@ -51,10 +51,11 @@ public enum Role {
* 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}, {@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}).
* {@code ASK} only as its own pane, {@code SEND} only to a target that would itself resolve
* as a lead ({@link #PRIMARY}) or as {@code OBSERVER}, and {@code TASK_READ} only a ticket its
* own {@code fleet_send(wait:false)} created; may not {@code SPAWN}/{@code STOP}/
* {@code DRAIN}/{@code HANDOVER}, read a session's status or another caller's ticket, or reach
* the coordination broker ({@code COORD_SEND}/{@code COORD_READ}).
*/
OBSERVER,
@@ -137,6 +137,18 @@ public final class AgentControl {
return result.path("read").path("text").asText("");
}
/**
* Read an agent's terminal with its ANSI styling kept, instead of the stripped text {@link
* #read} returns. Needed when a caller must tell apart text the pane draws dim (a placeholder
* hint) from text drawn plain (the operator's own typing).
*
* @param source one of {@code visible|recent|recent_unwrapped|detection}
*/
public String readWithStyling(String target, String source) {
JsonNode result = agentCall("agent.read", target, Map.of("source", source, "strip_ansi", false));
return result.path("read").path("text").asText("");
}
/** Current agent record (status, session UUID, pane). */
public Agent get(String target) {
return Agent.from(agentCall("agent.get", target, Map.of()).get("agent"));
@@ -3,9 +3,12 @@ package dev.ltms.fleet.herdr;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
/**
* Whether an agent pane's input box is clear for a delivery.
@@ -29,20 +32,22 @@ 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.
* herdr {@code agent.read} source, read with its ANSI styling kept. The input box is always
* drawn here, carrying transcript scrollback above it, which is why only the last box line is
* the live one. Styling must survive the read because the pane draws a placeholder hint — the
* pane's own last submitted prompt — in the same spot as unsubmitted text, dimmed; only the
* escape codes tell the two apart.
*/
static final String PROBE_SOURCE = "detection";
static final String PROBE_SOURCE = "visible";
/** 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.
* Input box markers, each matched only as a line's first characters once any leading ANSI
* escape codes are skipped: 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("❯", "│ >");
@@ -52,6 +57,15 @@ public final class PromptBox {
/** Block glyphs a terminal capture can leave in an otherwise empty box for the cursor cell. */
private static final String CURSOR_GLYPHS = "█▉▊▋▌▍▎▏";
/** An SGR escape sequence, e.g. {@code ESC[2m} (faint) or {@code ESC[0m} (reset). */
private static final Pattern SGR = Pattern.compile("\u001b\\[([0-9;]*)m");
/** The SGR code that dims text — herdr's placeholder hint is drawn inside a span of this. */
private static final String FAINT_CODE = "2";
/** The SGR code (or an empty code list) that clears every attribute, including faint. */
private static final List<String> RESET_CODES = List.of("", "0");
/** What a box holds: nothing, unsubmitted characters, or a pane this cannot read as a box. */
public enum State { EMPTY, DRAFT, UNREADABLE }
@@ -94,7 +108,7 @@ public final class PromptBox {
private Reading inspect(String target) {
String pane;
try {
pane = agents.read(target, PROBE_SOURCE);
pane = agents.readWithStyling(target, PROBE_SOURCE);
} catch (RuntimeException e) {
log.debug("prompt box read for {} failed, holding delivery: {}", target, e.getMessage());
return new Reading(State.UNREADABLE, 0);
@@ -114,10 +128,9 @@ public final class PromptBox {
* 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.
* <p>Whitespace, a trailing box border and a cursor block count as nothing. A placeholder hint —
* text the pane draws faint, in the same spot as unsubmitted text — also counts as nothing: only
* a character drawn outside a faint span is the operator's own typing.
*/
static Reading classify(String pane) {
if (pane == null || pane.isBlank()) return new Reading(State.UNREADABLE, 0);
@@ -142,28 +155,77 @@ public final class PromptBox {
return found;
}
/** Length of the box marker this line starts with, or {@code 0} if it starts with none. */
/**
* Length of the marker prefix — any leading SGR escape codes, then a box marker — this line
* starts with, or {@code 0} if it starts with neither. The colour drawn on the caret itself
* (e.g. an empty box's grey) sits before the marker glyph, so it must be skipped before the
* marker can match.
*/
private static int markerLength(String line) {
int skip = leadingEscapeLength(line);
for (String marker : BOX_MARKERS) {
if (line.startsWith(marker)) return marker.length();
if (line.startsWith(marker, skip)) return skip + marker.length();
}
return 0;
}
/** Length of the run of SGR escape codes starting at the beginning of {@code line}. */
private static int leadingEscapeLength(String line) {
Matcher m = SGR.matcher(line);
int pos = 0;
while (m.find(pos) && m.start() == pos) pos = m.end();
return pos;
}
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. */
/** One rendered character of a box line, and whether it was drawn inside a faint (dim) span. */
private record Glyph(char c, boolean faint) {
}
/**
* The text the box holds: its own line after the marker, with border, padding, cursor and any
* faint (placeholder-hint) text left out — only a character drawn outside a faint span is the
* operator's own typing.
*/
private static String boxContent(String boxLine) {
String line = boxLine.substring(markerLength(boxLine)).stripTrailing();
if (line.endsWith("│")) line = line.substring(0, line.length() - 1);
List<Glyph> glyphs = renderedGlyphs(boxLine.substring(markerLength(boxLine)));
int end = glyphs.size();
while (end > 0 && isBoxPadding(glyphs.get(end - 1).c())) end--;
if (end > 0 && glyphs.get(end - 1).c() == '│') end--;
StringBuilder content = new StringBuilder();
for (char c : line.toCharArray()) {
if (Character.isWhitespace(c) || CURSOR_GLYPHS.indexOf(c) >= 0) continue;
content.append(c);
for (int i = 0; i < end; i++) {
Glyph glyph = glyphs.get(i);
if (glyph.faint() || isBoxPadding(glyph.c())) continue;
content.append(glyph.c());
}
return content.toString();
}
/** Decode {@code text} into its rendered characters, tracking the faint (SGR 2) span each sits in. */
private static List<Glyph> renderedGlyphs(String text) {
List<Glyph> glyphs = new ArrayList<>();
Matcher m = SGR.matcher(text);
boolean faint = false;
int i = 0;
while (i < text.length()) {
if (m.find(i) && m.start() == i) {
String codes = m.group(1);
if (RESET_CODES.contains(codes)) faint = false;
else if (FAINT_CODE.equals(codes)) faint = true;
i = m.end();
continue;
}
glyphs.add(new Glyph(text.charAt(i), faint));
i++;
}
return glyphs;
}
private static boolean isBoxPadding(char c) {
return Character.isWhitespace(c) || Character.isSpaceChar(c) || CURSOR_GLYPHS.indexOf(c) >= 0;
}
}
@@ -132,6 +132,12 @@ public final class Injector {
* already uses for the same purpose.
*/
private final LongSupplier nowMillis;
/**
* Mail offered to panes that collect it themselves. Owned here because this is the single
* writer of delivery state, and the offer must be made and taken back under the same target
* monitor that guards the queue the message is still sitting on.
*/
private final PaneInbox paneInbox;
private final ConcurrentHashMap<String, Target> targets = new ConcurrentHashMap<>();
/** Delivery only; completion signalling is a no-op and every target is treated as available. */
@@ -203,6 +209,7 @@ public final class Injector {
this.ready = ready;
this.forget = forget;
this.nowMillis = nowMillis;
this.paneInbox = new PaneInbox(nowMillis);
}
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
@@ -232,6 +239,7 @@ public final class Injector {
this.ready = ready;
this.forget = forget;
this.nowMillis = nowMillis;
this.paneInbox = new PaneInbox(nowMillis);
}
private AgentControl agentsFor(String target) {
@@ -348,6 +356,10 @@ public final class Injector {
boolean awaitingPostTurnPickup;
boolean postTurnObserved;
int injectableSincePostTurnPickup;
/** The head of {@link #queue} as offered to a mod-served pane, or {@code null}. */
PaneInbox.Entry inboxOffer;
/** Whether the delivery the pickup latch is waiting on was collected rather than typed. */
boolean deliveredViaInbox;
synchronized void add(Pending p) {
queue.add(p);
@@ -376,7 +388,8 @@ public final class Injector {
/**
* Cancel this exact queued delivery. The target monitor serializes this operation with
* {@link #onStatus}: if delivery wins that race, this returns {@link Cancellation#DELIVERED}
* rather than claiming the message remained queued.
* rather than claiming the message remained queued. A message a mod-served pane has already
* collected answers the same way, even though no poll has recorded that delivery yet.
*/
public Cancellation cancel(Delivery delivery) {
Pending p = delivery.pending;
@@ -385,7 +398,22 @@ public final class Injector {
return cancellationOf(p);
}
synchronized (t) {
if (p.state != Pending.State.QUEUED || !t.queue.remove(p)) {
if (p.state != Pending.State.QUEUED) {
return cancellationOf(p);
}
if (t.queue.peek() == p && t.inboxOffer != null) {
// This exact entry is the one offered to a mod-served pane. The pane takes an
// offer on its own thread, so withdraw first and then read the outcome: a taken
// offer means the pane already holds this text, and the next poll records that
// delivery. Cancelling it would tell the caller nothing arrived while the pane
// acts on it.
paneInbox.withdrawAll(p.target);
if (t.inboxOffer.taken()) {
return Cancellation.DELIVERED;
}
t.inboxOffer = null;
}
if (!t.queue.remove(p)) {
return cancellationOf(p);
}
p.state = Pending.State.CANCELLED;
@@ -470,10 +498,14 @@ public final class Injector {
t.injectableSincePickup = 0;
t.awaitingCompletion = false;
t.turnObserved = false;
} else {
} else if (!t.deliveredViaInbox) {
// Delivered but still idle → the worker hasn't picked it up; the submit
// keystroke likely raced the paste (esp. right as the TUI became ready).
// Re-nudge Enter (CB-113) until the worker starts (WORKING) or the grace ends.
//
// A pane that collected the message submits it itself, and nothing was
// typed into it. Pressing Enter there would submit whatever its operator
// has in the prompt box instead.
resubmit = true;
}
}
@@ -498,45 +530,77 @@ public final class Injector {
if (p != null && ready.test(target)) {
t.notReadySincePoll = 0;
t.notReadySinceMillis = 0;
// fleetd #551: poll and record BEFORE the irreversible send, not after.
// The entry comes off the queue and its state is set to ATTEMPTED here,
// unconditionally — so a Throwable escaping the send call below (caught or
// not) can never leave the entry QUEUED at the head of t.queue (the fleetd
// #546 hazard, since peek() alone would let the next onStatus round re-enter
// this block and send the same text again), and no path can write a
// confident DELIVERED or NOT_DELIVERED before we actually know which one
// happened.
t.queue.poll();
p.state = Pending.State.ATTEMPTED;
try {
agentsFor(target).send(target, p.text());
p.state = Pending.State.DELIVERED;
t.awaitingPickup = true;
t.awaitingCompletion = true;
t.turnObserved = false;
t.injectableSincePickup = 0;
sent = p;
} catch (Throwable e) {
// fleetd #551: leave p.state == ATTEMPTED (recorded above, before the
// call) rather than downgrading it to NOT_DELIVERED here — reaching this
// catch does not prove the text never reached the pane. Three of the
// four HerdrException throw sites in HerdrCodec fire only after herdr
// has already replied (so it processed the request), and the fourth (a
// transport IOException) leaves it genuinely unknown whether herdr even
// received the bytes — see #551 comment 16867. The one exception is a
// herdr `*_not_found` error: that family is already read as "definitely
// absent, not merely inconclusive" everywhere else in this codebase
// (StatusPoller, AgentControl's own retry, WorkspaceControl,
// HerdrPeerLauncher, FleetApp, ReplyPushLoop) because it means the
// target pane/agent does not exist at all, so nothing could have been
// pasted anywhere — #551 keeps the new state consistent with that
// existing vocabulary rather than inventing a second one.
if (e instanceof HerdrException he && he.code() != null
&& he.code().endsWith("_not_found")) {
p.state = Pending.State.NOT_DELIVERED;
boolean modServed = paneInbox.isModServed(target);
if (t.inboxOffer != null && !modServed) {
// The pane stopped collecting its mail, so take the offer back and
// fall through to the terminal route below. withdrawAll leaves an
// entry the pane took first alone, so the branch under it still sees
// that as the delivery it is.
paneInbox.withdrawAll(target);
if (!t.inboxOffer.taken()) {
t.inboxOffer = null;
}
}
if (t.inboxOffer == null && modServed) {
t.inboxOffer = paneInbox.offer(target, p.text());
}
if (t.inboxOffer != null) {
// The message stays at the head of the queue until the pane takes
// it: nothing has reached the pane yet, so nothing may be recorded
// as delivered and nothing may be failed.
if (t.inboxOffer.taken()) {
t.inboxOffer = null;
t.queue.poll();
p.state = Pending.State.DELIVERED;
t.awaitingPickup = true;
t.awaitingCompletion = true;
t.turnObserved = false;
t.injectableSincePickup = 0;
t.deliveredViaInbox = true;
sent = p;
}
} else {
// fleetd #551: poll and record BEFORE the irreversible send, not after.
// The entry comes off the queue and its state is set to ATTEMPTED here,
// unconditionally — so a Throwable escaping the send call below (caught or
// not) can never leave the entry QUEUED at the head of t.queue (the fleetd
// #546 hazard, since peek() alone would let the next onStatus round re-enter
// this block and send the same text again), and no path can write a
// confident DELIVERED or NOT_DELIVERED before we actually know which one
// happened.
t.queue.poll();
p.state = Pending.State.ATTEMPTED;
try {
agentsFor(target).send(target, p.text());
p.state = Pending.State.DELIVERED;
t.awaitingPickup = true;
t.awaitingCompletion = true;
t.turnObserved = false;
t.injectableSincePickup = 0;
t.deliveredViaInbox = false;
sent = p;
} catch (Throwable e) {
// fleetd #551: leave p.state == ATTEMPTED (recorded above, before the
// call) rather than downgrading it to NOT_DELIVERED here — reaching this
// catch does not prove the text never reached the pane. Three of the
// four HerdrException throw sites in HerdrCodec fire only after herdr
// has already replied (so it processed the request), and the fourth (a
// transport IOException) leaves it genuinely unknown whether herdr even
// received the bytes — see #551 comment 16867. The one exception is a
// herdr `*_not_found` error: that family is already read as "definitely
// absent, not merely inconclusive" everywhere else in this codebase
// (StatusPoller, AgentControl's own retry, WorkspaceControl,
// HerdrPeerLauncher, FleetApp, ReplyPushLoop) because it means the
// target pane/agent does not exist at all, so nothing could have been
// pasted anywhere — #551 keeps the new state consistent with that
// existing vocabulary rather than inventing a second one.
if (e instanceof HerdrException he && he.code() != null
&& he.code().endsWith("_not_found")) {
p.state = Pending.State.NOT_DELIVERED;
}
sent = p;
sendError = e;
}
sent = p;
sendError = e;
}
} else if (p != null) {
// fleetd #501: stamp the wall-clock time of the FIRST non-ready sample in
@@ -555,6 +619,8 @@ public final class Injector {
for (Pending pending : notReady) {
pending.state = Pending.State.NOT_DELIVERED;
}
paneInbox.withdrawAll(target);
t.inboxOffer = null;
// fleetd #501: t.notReadySincePoll — the loop's own counter, already in
// scope — is printed here instead of the READINESS_GRACE_POLLS constant.
// On this branch the counter has JUST reached the threshold, so the two
@@ -633,6 +699,8 @@ public final class Injector {
for (Pending pending : queueStalled) {
pending.state = Pending.State.NOT_DELIVERED;
}
paneInbox.withdrawAll(target);
t.inboxOffer = null;
log.warn("queue for {} never drained after {} polls (limit={} polls/{}s): "
+ "failing {} queued message(s) that were never attempted",
target, t.queueWaitSincePoll, QUEUE_WAIT_GRACE_POLLS,
@@ -846,6 +914,24 @@ public final class Injector {
}
}
/**
* Hand {@code terminal} every message held for it and stamp it as collecting its own mail.
* While that stamp is fresh this injector offers that pane's messages instead of typing them;
* once it goes stale the pane's queued mail takes the terminal route again.
*
* <p>Returns the messages in the order they were queued, and an empty list when there are
* none — an empty collection still counts as collecting, so a pane that polls on a timer stays
* mod-served between messages.
*/
public List<String> collectInbox(String terminal) {
return paneInbox.drain(terminal);
}
/** Whether {@code terminal} has collected its mail recently enough to be offered the next one. */
public boolean isModServed(String terminal) {
return paneInbox.isModServed(terminal);
}
/**
* Targets the poller must keep sampling: those with a queued message, an awaited pickup, or an
* awaited turn completion (so the {@code working → idle} boundary is observed).
@@ -899,6 +985,11 @@ public final class Injector {
p.state = Pending.State.NOT_DELIVERED;
}
t.queue.clear();
// The pane is gone, so drop its offered mail and its poll stamp together: a terminal
// id can be reused, and a stale stamp would make the next pane under it look mod-served
// before it has ever collected anything.
paneInbox.forget(target);
t.inboxOffer = null;
hadDeliveredTurn = t.awaitingCompletion;
t.awaitingCompletion = false;
t.awaitingPickup = false;
@@ -0,0 +1,141 @@
package dev.ltms.fleet.inject;
import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Deque;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.LongSupplier;
/**
* Mail held for a pane that collects it itself instead of having it typed into its terminal.
*
* <p>A pane becomes <em>mod-served</em> by calling {@code fleet_inbox}: {@link #drain} stamps the
* pane as polling, and {@link #isModServed} answers {@code true} while that stamp is younger than
* {@link #MOD_SERVED_WINDOW_MILLIS}. Nothing else sets it, so a pane that has never polled is
* never mod-served and its mail takes the terminal route.
*
* <p>An offered entry is removed exactly once, by {@link #drain} or by {@link #withdrawAll}, and
* both run under the owning pane's monitor. So an entry the pane took is never also withdrawn, and
* an entry that was withdrawn can never still be collected — which is what lets the {@link
* Injector} keep one message both offered here and queued for the terminal without risking two
* deliveries of it.
*
* <p>This class holds no queue of its own beyond what is currently offered: the {@link Injector}
* keeps the message on its own queue until the pane takes it, so a pane that stops polling strands
* nothing.
*/
public class PaneInbox {
/**
* How long after a {@link #drain} a pane still counts as mod-served. It must exceed the mod's
* own poll interval by enough that a few missed polls are not read as a pane that stopped,
* while staying short enough that a pane which really stopped falls back to the terminal route
* promptly.
*/
public static final long MOD_SERVED_WINDOW_MILLIS = 15_000;
/** One message held for a pane until that pane collects it. */
public static final class Entry {
private final String text;
private boolean taken; // written under the owning pane's monitor
private Entry(String text) {
this.text = text;
}
/** The message text, as it will be handed to the pane. */
public String text() {
return text;
}
/** Whether the pane has collected this entry. Once {@code true} it never goes back. */
public synchronized boolean taken() {
return taken;
}
private synchronized void markTaken() {
taken = true;
}
}
private final LongSupplier nowMillis;
private final ConcurrentHashMap<String, Long> lastPolledAtMillis = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Deque<Entry>> offered = new ConcurrentHashMap<>();
public PaneInbox() {
this(System::currentTimeMillis);
}
public PaneInbox(LongSupplier nowMillis) {
this.nowMillis = Objects.requireNonNull(nowMillis, "nowMillis");
}
/**
* Whether {@code terminal} collected its mail within {@link #MOD_SERVED_WINDOW_MILLIS}. A
* terminal that has never collected any is never mod-served.
*/
public boolean isModServed(String terminal) {
Long at = terminal == null ? null : lastPolledAtMillis.get(terminal);
return at != null && nowMillis.getAsLong() - at <= MOD_SERVED_WINDOW_MILLIS;
}
/** Hold {@code text} for {@code terminal} to collect, and return the entry holding it. */
public Entry offer(String terminal, String text) {
Entry entry = new Entry(text);
Deque<Entry> queue = offered.computeIfAbsent(terminal, _ -> new ArrayDeque<>());
synchronized (queue) {
queue.add(entry);
}
return entry;
}
/**
* Take back every entry {@code terminal} has not collected. An entry the pane took first stays
* taken — this never un-delivers one.
*/
public void withdrawAll(String terminal) {
Deque<Entry> queue = terminal == null ? null : offered.get(terminal);
if (queue == null) {
return;
}
synchronized (queue) {
queue.removeIf(e -> !e.taken());
}
}
/**
* Collect everything held for {@code terminal}, in the order it was offered, and stamp the pane
* as polling. Each returned entry is marked taken, so the {@link Injector} can tell a message
* the pane really has from one it merely offered.
*/
public List<String> drain(String terminal) {
if (terminal == null || terminal.isBlank()) {
return List.of();
}
lastPolledAtMillis.put(terminal, nowMillis.getAsLong());
Deque<Entry> queue = offered.get(terminal);
if (queue == null) {
return List.of();
}
List<String> collected = new ArrayList<>();
synchronized (queue) {
for (Entry e : queue) {
e.markTaken();
collected.add(e.text());
}
queue.clear();
}
return collected;
}
/** Forget a pane that is gone, so neither its poll stamp nor its offered mail lingers. */
public void forget(String terminal) {
if (terminal == null) {
return;
}
lastPolledAtMillis.remove(terminal);
offered.remove(terminal);
}
}
@@ -121,9 +121,9 @@ public final class FleetMcp {
/**
* Kept as a field (rather than only captured by the {@code contextExtractor} closure) so
* {@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.
* CallerResolver#observerSendTarget()} — 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
@@ -323,12 +323,12 @@ public final class FleetMcp {
* 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}
* @param observerSendTarget the same predicate {@link CallerResolver#observerSendTarget}
* 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) {
Predicate<String> deliverable, Predicate<String> architectSlot, Predicate<String> observerSendTarget) {
/** 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);
@@ -528,8 +528,13 @@ public final class FleetMcp {
// name would never match there anyway, but passing null for them keeps the intent
// explicit rather than relying on that lookup to filter it out.
String delegatorName = caller.isPrimary() ? caller.name() : null;
Runnable onAccepted = () ->
primaryRegistry.recordDelegation(target, callerTerminal, delegatorName);
// A nudge about this target must never tell its recipient to run a call Authz
// would refuse it — asked once, here, while the real Principal is still in
// scope, rather than re-derived from a bare terminal string later.
boolean mayDrainNudge = Authz.permits(caller, Authz.Action.DRAIN, target);
boolean mayAnswerNudge = Authz.permits(caller, Authz.Action.ANSWER, target);
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, callerTerminal,
delegatorName, mayDrainNudge, mayAnswerNudge);
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
@@ -553,6 +558,16 @@ public final class FleetMcp {
if (denied != null) return denied;
return ask(messages, self, str(req.arguments(), "question"), timeoutMs(req.arguments()));
};
// fleet_inbox: the caller collects the mail queued for its OWN pane. Identity is the
// CONNECTION, never an argument, and the authz check is terminal ownership — the same
// shape as fleet_reply above.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> inboxHandler =
(exchange, req) -> {
String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_inbox", req.arguments()), self);
if (denied != null) return denied;
return inbox(messages, self);
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> statusHandler =
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null);
@@ -564,10 +579,18 @@ public final class FleetMcp {
Map<String, Object> a = req.arguments();
String target = str(a, "target");
String coordId = str(a, "coordId");
String ticket = str(a, "ticket");
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
Authz.Action action = toolAction("fleet_poll", a);
// TASK_READ reads a ticket, not a DRAIN target -- the gate must see the
// ticket id, and an observer's grant is confined to a ticket it created.
String authzTarget = action == Authz.Action.TASK_READ ? ticket : target;
Predicate<String> observerOwnsTicket = action == Authz.Action.TASK_READ
? t -> messages.ownsTicket(t, principal(exchange).ownerKey())
: Authz.NO_OBSERVER_OWNED_TICKET;
McpSchema.CallToolResult denied = deny(exchange, action, authzTarget, observerOwnsTicket);
if (denied != null) return denied;
return poll(messages, leadChannel, str(a, "ticket"), target, coordId,
return poll(messages, leadChannel, ticket, target, coordId,
principal(exchange).ownerKey());
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
@@ -606,7 +629,7 @@ public final class FleetMcp {
PaneSource panes = new PaneSource(() -> identity.panes().tabLabelsByTabId(),
() -> identity.panes().workspaceLabelsByWorkspaceId(),
Fleetd.deliverableTo(presence, callers::leads, callers::collaborators),
callers::boundToArchitectSlot, callers.sendableObserverTarget());
callers::boundToArchitectSlot, callers.observerSendTarget());
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
callerTerminal(exchange),
@@ -658,6 +681,7 @@ public final class FleetMcp {
McpSchema.Tool fleetProfiles = profilesTool();
McpSchema.Tool fleetWhoami = whoamiTool();
McpSchema.Tool fleetHandover = handoverTool();
McpSchema.Tool fleetInbox = inboxTool();
// fleetd #469: the tool schemas above are already named from FleetTool.wireName(), but
// this is the check that a schema was not accidentally dropped, duplicated, or added
@@ -668,7 +692,7 @@ public final class FleetMcp {
Set<String> registeredToolNames = Set.of(fleetSend.name(), fleetReply.name(), fleetAsk.name(),
fleetStatus.name(), fleetPoll.name(), fleetAck.name(), fleetSpawn.name(),
fleetList.name(), fleetStop.name(), fleetProfiles.name(), fleetWhoami.name(),
fleetHandover.name());
fleetHandover.name(), fleetInbox.name());
if (!registeredToolNames.equals(FleetTool.wireNames())) {
throw new IllegalStateException("fleetd #469: registered MCP tools " + registeredToolNames
+ " do not match the canonical tool set " + FleetTool.wireNames()
@@ -690,6 +714,7 @@ public final class FleetMcp {
.toolCall(fleetProfiles, profilesHandler)
.toolCall(fleetWhoami, whoamiHandler)
.toolCall(fleetHandover, handoverHandler)
.toolCall(fleetInbox, inboxHandler)
.build();
this.metrics = metrics;
}
@@ -733,6 +758,16 @@ public final class FleetMcp {
return denyFor(principal(exchange), action, target);
}
/**
* As {@link #deny(McpSyncServerExchange, Authz.Action, String)}, with a real classifier for
* an observer's {@code TASK_READ} — the only call site that can supply one is a ticket poll,
* which knows the ticket id and the caller's owner key.
*/
private McpSchema.CallToolResult deny(McpSyncServerExchange exchange, Authz.Action action,
String target, Predicate<String> observerOwnsTicket) {
return denyFor(principal(exchange), action, target, observerOwnsTicket);
}
/**
* The policy half of {@link #deny}: everything except pulling the caller out of the MCP
* exchange. Kept separate so the authorization decision — the actual control — is unit-testable
@@ -746,13 +781,22 @@ public final class FleetMcp {
* @return {@code null} when the call may proceed, or the error result to return when it may not
*/
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target) {
return denyFor(caller, action, target, Authz.NO_OBSERVER_OWNED_TICKET);
}
/**
* As {@link #denyFor(Principal, Authz.Action, String)}, with a real classifier for an
* observer's {@code TASK_READ}.
*/
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target,
Predicate<String> observerOwnsTicket) {
// The enforcement switch lives HERE rather than in the exchange-facing wrapper: any future
// tool that calls this directly must not be able to skip the gate by accident.
if (!authorizationEnforced) {
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
}
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator(),
callers.sendableObserverTarget())) {
callers.observerSendTarget(), observerOwnsTicket)) {
if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
}
@@ -814,13 +858,15 @@ public final class FleetMcp {
}
/**
* Who may see {@code fleet_list}'s {@code leads} array — exactly the roles that may
* {@link Authz.Action#SEND} to a lead: the primary, an architect, and a collaborator. A
* collaborator's own {@code fleet_whoami} carries no lead address, and {@code leads} is the
* only place this tool gives one, so a collaborator needs this array to use the send it
* already holds. A worker can never {@code SEND} at all, so it still sees neither this array
* nor {@code members}; a worker's own facts come from {@code fleet_whoami} instead. Split
* out for the same reason as {@link #coordinatorVisibleTo} and
* Who may see {@code fleet_list}'s {@code leads} array — the primary, an architect, and a
* collaborator. A collaborator's own {@code fleet_whoami} carries no lead address, and this is
* the only place this tool gives one, so a collaborator needs this array to use the
* {@link Authz.Action#SEND} it already holds. An observer holds that send to a lead too, but
* learns the address from its filtered {@code panes} rows instead: a {@code leads} row carries
* a lead's name, its configured context window and its config dir, which are the fleet's own
* shape rather than an address. A worker can never {@code SEND} at all, so it still sees
* neither this array nor {@code members}; a worker's own facts come from {@code fleet_whoami}
* instead. Split out for the same reason as {@link #coordinatorVisibleTo} and
* {@link #collaboratorsVisibleTo}: the decision must be unit-testable without fabricating an
* SDK {@code McpSyncServerExchange}, and the handler must call this named predicate rather
* than inlining the check.
@@ -843,9 +889,10 @@ public final class FleetMcp {
/**
* 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}.
* {@code SEND}, so it still does not see this array. An observer does hold {@code SEND}, to a
* lead or another observer pane, so it sees the array too — and it is the only place this tool
* gives it a lead's address. {@code listFleet} filters its rows to
* {@link CallerResolver#observerSendTarget} and reduces each one; see {@code paneRows}.
*/
static boolean panesVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator() || caller.isObserver();
@@ -1286,6 +1333,7 @@ public final class FleetMcp {
case SEND -> sendAction(str(arguments, "coordId"), str(arguments, "turnId"));
case REPLY -> Authz.Action.REPLY;
case ASK -> Authz.Action.ASK;
case INBOX -> Authz.Action.INBOX;
case STATUS -> Authz.Action.TASK_READ;
case LIST, PROFILES, WHOAMI -> Authz.Action.READ;
case POLL -> pollAction(str(arguments, "target"), str(arguments, "coordId"));
@@ -1429,6 +1477,33 @@ public final class FleetMcp {
return text(outcome.description());
}
/**
* {@code fleet_inbox}: hand the caller the messages queued for its own pane, so it submits them
* itself instead of having them typed into its terminal. {@code callerTerminal} is resolved
* from the connection; there is deliberately no pane argument, so no caller can collect another
* pane's mail.
*
* <p>Calling this is also what marks the pane as collecting its own mail, for a window the
* injector re-checks before every delivery. A pane that stops calling it stops being served
* this way and its queued mail is typed instead, so an empty answer is a normal result that
* must still be requested on a timer.
*
* <p>Returns a JSON object with {@code count} and {@code messages} (in the order they were
* queued), so a caller can tell "no mail" apart from a failure.
*/
static McpSchema.CallToolResult inbox(MessageService messages, String callerTerminal) {
if (callerTerminal == null) {
return error("fleet_inbox could not identify the calling pane from the connection, so "
+ "there is no inbox to collect");
}
List<String> collected = messages.collectInbox(callerTerminal);
Map<String, Object> m = new LinkedHashMap<>();
m.put("sessionId", callerTerminal);
m.put("count", collected.size());
m.put("messages", collected);
return text(json(m));
}
/** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */
static McpSchema.CallToolResult ack(MessageService messages, String target, String msgId) {
if (isBlank(target) || isBlank(msgId)) {
@@ -2115,8 +2190,8 @@ public final class FleetMcp {
}
/**
* 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
* As above, plus: an observer sees the {@code panes} array too, but filtered to
* {@link CallerResolver#observerSendTarget} 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
@@ -2539,9 +2614,10 @@ public final class FleetMcp {
* 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}.
* <p>For an observer caller ({@code observerView}), the rows are filtered to
* {@link PaneSource#observerSendTarget} before being built, and each row is reduced — see
* {@code paneRow}. A lead's pane survives that filter, so it is where an observer reads a
* lead's {@code sessionId}.
*/
private static List<Map<String, Object>> paneRows(Map<String, Agent> live, List<MemberSession> roster,
Map<String, String> leads, Map<String, String> collaborators, PaneSource panes,
@@ -2552,7 +2628,7 @@ public final class FleetMcp {
.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()))
.filter(a -> !observerView || panes.observerSendTarget().test(a.terminalId()))
.sorted(Comparator.comparing(Agent::terminalId))
.map(a -> paneRow(a, byTerminal.get(a.terminalId()), leads, collaborators, tabLabels,
workspaceLabels, panes, observerView))
@@ -2567,9 +2643,9 @@ public final class FleetMcp {
* @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},
* @param observerView 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)
*/
@@ -2817,9 +2893,12 @@ public final class FleetMcp {
return tool(FleetTool.LIST.wireName(),
"List the whole fleet the bridge tracks, in two parts. 'members' is visible to "
+ "the primary and an architect only. 'leads' is visible to those two AND a "
+ "collaborator — exactly the roles that may fleet_send to a lead, so a "
+ "collaborator can learn a lead's sessionId before using the send it already "
+ "holds. A worker holds READ to call this tool at all, but gets neither "
+ "collaborator, so a collaborator can learn a lead's sessionId before using "
+ "the send it already holds. An observer may fleet_send to a lead too, but "
+ "reads that sessionId from its own 'panes' rows instead: a 'panes' row for "
+ "an observer is filtered to the panes it may send to — a lead's pane and "
+ "another observer's — and reduced to sessionId, label, status, role and "
+ "deliverable. A worker holds READ to call this tool at all, but gets neither "
+ "array, never an empty one; a worker reads its own session, "
+ "profile, state, worktree, branch and owner from fleet_whoami instead. "
+ "'leads' are your PEERS — other "
@@ -2905,6 +2984,18 @@ public final class FleetMcp {
objectSchema(Map.of(), List.of()));
}
private static McpSchema.Tool inboxTool() {
return tool(FleetTool.INBOX.wireName(),
"Collect the messages the fleet has queued for YOUR OWN pane, then act on each one. "
+ "Takes no arguments: the pane is resolved from your connection, so you can "
+ "never read another session's mail. Any role may call it for itself. Call it "
+ "on a timer — each call is also what tells the daemon you collect your own "
+ "mail, so it offers your next message here instead of typing it into your "
+ "terminal; stop calling it and your mail is typed instead. An empty "
+ "'messages' array is the normal answer when nothing is waiting.",
objectSchema(Map.of(), List.of()));
}
private static McpSchema.Tool handoverTool() {
return tool(FleetTool.HANDOVER.wireName(),
"Replace your OWN lead session once its context is full: write a handover file, "
@@ -44,7 +44,8 @@ public enum FleetTool {
STOP("fleet_stop"),
PROFILES("fleet_profiles"),
WHOAMI("fleet_whoami"),
HANDOVER("fleet_handover");
HANDOVER("fleet_handover"),
INBOX("fleet_inbox");
private final String wireName;
@@ -48,8 +48,13 @@ public final class PrimaryRegistry {
*/
private final ConcurrentHashMap<String, Delegation> leadByTarget = new ConcurrentHashMap<>();
/** A recorded delegator: the terminal learned from call traffic, and its name, if it has one. */
private record Delegation(String terminal, String name) {
/**
* A recorded delegator: the terminal learned from call traffic, its name if it has one, and
* whether a nudge may tell it to {@code fleet_poll(target=…)} or {@code fleet_send(turnId=…)}
* — the real authorization decision for this delegator's role, captured once by the caller at
* delegation time rather than re-derived from a bare terminal string later.
*/
private record Delegation(String terminal, String name, boolean mayDrainNudge, boolean mayAnswerNudge) {
}
/**
@@ -128,14 +133,28 @@ public final class PrimaryRegistry {
/**
* As {@link #recordDelegation(String, String)}, additionally recording the delegating lead's
* name when the caller carries one. See {@link #record(String, String)} for why the name
* matters.
* name when the caller carries one, and granting a nudge about this target everything a
* primary may be told to run. See {@link #record(String, String)} for why the name matters,
* and {@link #recordDelegation(String, String, String, boolean, boolean)} for why every other
* caller must supply its own, real grant instead of this default.
*/
public void recordDelegation(String target, String leadTerminal, String leadName) {
recordDelegation(target, leadTerminal, leadName, true, true);
}
/**
* As {@link #recordDelegation(String, String, String)}, with the delegating caller's own
* {@code DRAIN} and {@code ANSWER} grants — the real {@code Authz.permits} decision for its
* role, asked once by the caller at delegation time, since a terminal string alone cannot be
* resolved back to a role here. A nudge about this target must offer only what
* {@link #nudgeMayDrainFor(String)}/{@link #nudgeMayAnswerFor(String)} report back as true.
*/
public void recordDelegation(String target, String leadTerminal, String leadName,
boolean mayDrainNudge, boolean mayAnswerNudge) {
if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) {
return;
}
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName)));
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName), mayDrainNudge, mayAnswerNudge));
}
/** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */
@@ -167,6 +186,32 @@ public final class PrimaryRegistry {
return currentPrimaryTerminal();
}
/**
* Whether a nudge about {@code target} may tell its recipient to {@code fleet_poll(target=…)}
* — present on exactly the same condition as {@link #nudgeTargetFor(String)}, carrying the
* grant recorded for a delegation, or {@code true} for the singleton primary fallback, which
* is always genuinely the primary.
*/
public Optional<Boolean> nudgeMayDrainFor(String target) {
Delegation delegation = target == null ? null : leadByTarget.get(target);
if (delegation != null) {
return Optional.of(delegation.mayDrainNudge());
}
return currentPrimaryTerminal().isPresent() ? Optional.of(true) : Optional.empty();
}
/**
* As {@link #nudgeMayDrainFor(String)}, for whether a nudge may tell its recipient to
* {@code fleet_send(turnId=…)}.
*/
public Optional<Boolean> nudgeMayAnswerFor(String target) {
Delegation delegation = target == null ? null : leadByTarget.get(target);
if (delegation != null) {
return Optional.of(delegation.mayAnswerNudge());
}
return currentPrimaryTerminal().isPresent() ? Optional.of(true) : Optional.empty();
}
/**
* The known primary terminal, or empty if not yet learned (and not pinned) — the raw value as
* it was recorded, with no attempt to resolve a named lead's current pane. Callers that need a
@@ -514,6 +514,20 @@ public final class MessageService {
return false;
}
/**
* Hand {@code session} every message queued for it that it has not collected yet, and record
* that it collects its own mail. While that record is fresh, delivery to that session is
* offered for collection instead of typed into its terminal; once it goes stale, the terminal
* route takes over again with nothing lost.
*
* <p>The messages are returned in the order they were queued, and are removed by this call.
* An empty list is an ordinary answer: a session polling on a timer keeps itself collecting
* between messages.
*/
public List<String> collectInbox(String session) {
return injector.collectInbox(session);
}
/**
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
@@ -1478,6 +1492,16 @@ public final class MessageService {
|| Objects.equals(callerOwner, task.creatorOwner);
}
/**
* Whether {@code callerOwner} owns {@code ticket}, with no side effect — unlike
* {@link #poll(String, String)}, this never fires the ticket's collection hook. A ticket this
* daemon has never heard of is owned by nobody.
*/
public boolean ownsTicket(String ticket, String callerOwner) {
Task task = tasks.get(ticket);
return task != null && ownsTicket(task, callerOwner);
}
/**
* Test seam only — carries no production behaviour, and nothing in this class calls it;
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
@@ -178,13 +178,20 @@ public final class ReplyPushLoop {
pendingReplies.remove(target, owning);
continue;
}
// Never nudge a caller to run a DRAIN its own role could never perform -- the durable
// inbox stays the backstop for it, same as when no nudge target is known at all.
if (!owning.mayDrainNudge()) continue;
result.add(target);
}
return result;
}
/** A pending reply target: which lead to nudge, and how many nudges have named it so far. */
private record ReplyEntry(String lead, int nudgeCount) {
/**
* A pending reply target: which nudge target to notify, how many nudges have named it so
* far, and whether that nudge target's own role may actually run the
* {@code fleet_poll(target=…)} a reply nudge would tell it to run.
*/
private record ReplyEntry(String lead, int nudgeCount, boolean mayDrainNudge) {
}
/**
@@ -207,11 +214,12 @@ public final class ReplyPushLoop {
/**
* An open question awaiting the lead's answer: which ticket it belongs to, which worker asked,
* which lead to nudge, the question text, and how many nudges have named it so far (CB-598 —
* tracked per question, not per lead per source).
* which lead to nudge, the question text, how many nudges have named it so far (CB-598 —
* tracked per question, not per lead per source), and whether that nudge target's own role may
* actually run the {@code fleet_send(turnId=…)} a question nudge would tell it to run.
*/
private record PendingQuestion(String turnId, String ticket, String target, String lead,
String question, int nudgeCount) {
String question, int nudgeCount, boolean mayAnswerNudge) {
}
private record IncidentLead(String incidentId, String lead) {
@@ -229,7 +237,11 @@ public final class ReplyPushLoop {
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
private List<PendingQuestion> pendingQuestionsFor(String lead) {
return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList();
// Never nudge a caller to run an ANSWER its own role could never perform -- the 55-second
// ask window lapses on its own, same as when no nudge target is known at all.
return pendingQuestions.values().stream()
.filter(q -> lead.equals(q.lead()) && q.mayAnswerNudge())
.toList();
}
/** Question turnIds still open for {@code lead} — a plain snapshot for race comparison. */
@@ -385,7 +397,7 @@ public final class ReplyPushLoop {
boolean hasIncidentWork = !pendingIncidentKeysFor(lead).isEmpty();
boolean hasUnmappedTargetWork = !pendingUnmappedTargetKeysFor(lead).isEmpty();
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork && !hasUnmappedTargetWork) {
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
log.debug("push: nothing pending for nudge target {}, stopping reminder", lead);
return Action.STOP;
}
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
@@ -394,7 +406,7 @@ public final class ReplyPushLoop {
boolean incidentEligible = hasIncidentWork && incidentReminderCount < maxReminders;
boolean unmappedTargetEligible = hasUnmappedTargetWork && unmappedTargetReminderCount < maxReminders;
if (!replyEligible && !ticketEligible && !questionEligible && !incidentEligible && !unmappedTargetEligible) {
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
log.debug("push: reminder cap ({}) reached for nudge target {} on every source with pending work, stopping",
maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
@@ -403,15 +415,15 @@ public final class ReplyPushLoop {
try {
status = agents.status(lead);
} catch (RuntimeException e) {
log.debug("push: status check failed for lead {}, will retry", lead, e);
log.debug("push: status check failed for nudge target {}, will retry", lead, e);
return Action.WAIT_BUSY;
}
if (!status.injectable()) {
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
log.debug("push: nudge target {} 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);
log.debug("push: nudge target {} has unsubmitted text in its prompt box, waiting", lead);
return Action.WAIT_BUSY;
}
return Action.INJECT;
@@ -466,14 +478,14 @@ public final class ReplyPushLoop {
if (lead.isEmpty() || isLive(lead.get())) {
return lead;
}
log.debug("push: lead {} delegated to for {} is no longer live, forgetting the stale binding "
log.debug("push: nudge target {} delegated to for {} is no longer live, forgetting the stale binding "
+ "and falling back", lead.get(), target);
primaryRegistry.forgetDelegation(target);
Optional<String> fallback = primaryRegistry.nudgeTargetFor(target);
if (fallback.isEmpty() || isLive(fallback.get())) {
return fallback;
}
log.debug("push: fallback lead {} for {} is also not live, skipping this tick", fallback.get(), target);
log.debug("push: fallback nudge target {} for {} is also not live, skipping this tick", fallback.get(), target);
return Optional.empty();
}
@@ -496,9 +508,9 @@ public final class ReplyPushLoop {
} catch (RuntimeException e) {
boolean gone = e instanceof HerdrException he && "agent_not_found".equals(he.code());
if (gone) {
log.debug("push: lead {} no longer exists ({})", lead, e.toString());
log.debug("push: nudge target {} no longer exists ({})", lead, e.toString());
} else {
log.debug("push: liveness check for lead {} was inconclusive ({}); treating as live "
log.debug("push: liveness check for nudge target {} was inconclusive ({}); treating as live "
+ "rather than risk destroying a live binding", lead, e.toString());
}
return !gone;
@@ -517,11 +529,12 @@ public final class ReplyPushLoop {
public void onReplyQueued(String target) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
log.debug("push: no nudge target is known to be waiting on {}, skipping reminder", target);
return;
}
boolean mayDrainNudge = primaryRegistry.nudgeMayDrainFor(target).orElse(false);
pendingReplies.compute(target, (t, existing) ->
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount(), mayDrainNudge));
startOrCoalesce(lead.get());
}
@@ -545,7 +558,7 @@ public final class ReplyPushLoop {
public void onTicketTerminal(String ticket, String target, boolean failed) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
log.debug("push: no nudge target is known to be waiting on ticket {} (target {}), skipping nudge",
ticket, target);
return;
}
@@ -580,11 +593,13 @@ public final class ReplyPushLoop {
public void onQuestionOpened(String ticket, String target, String turnId, String question) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge",
log.debug("push: no nudge target is known to be waiting on {}'s question (turnId {}), skipping nudge",
target, turnId);
return;
}
pendingQuestions.put(turnId, new PendingQuestion(turnId, ticket, target, lead.get(), question, 0));
boolean mayAnswerNudge = primaryRegistry.nudgeMayAnswerFor(target).orElse(false);
pendingQuestions.put(turnId,
new PendingQuestion(turnId, ticket, target, lead.get(), question, 0, mayAnswerNudge));
startOrCoalesce(lead.get());
}
@@ -608,7 +623,7 @@ public final class ReplyPushLoop {
for (String target : targets) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.warn("push: backend incident {} has no known lead for target {}", incidentId, target);
log.warn("push: backend incident {} has no known nudge target for worker {}", incidentId, target);
continue;
}
targetsByLead.computeIfAbsent(lead.get(), _ -> new ArrayList<>()).add(target);
@@ -644,10 +659,10 @@ public final class ReplyPushLoop {
/** Start a reminder schedule for {@code lead}, or join the one already running. */
private void startOrCoalesce(String lead) {
if (activeLeads.putIfAbsent(lead, Boolean.TRUE) != null) {
log.debug("push: reminder loop already active for lead {}, work coalesced in", lead);
log.debug("push: reminder loop already active for nudge target {}, work coalesced in", lead);
return;
}
log.debug("push: starting reminder loop for lead {}", lead);
log.debug("push: starting reminder loop for nudge target {}", lead);
scheduleNext(lead);
}
@@ -752,11 +767,11 @@ public final class ReplyPushLoop {
|| pendingIncidentKeysFor(lead).stream().anyMatch(i -> !incidentsBefore.contains(i))
|| pendingUnmappedTargetKeysFor(lead).stream().anyMatch(i -> !unmappedTargetsBefore.contains(i));
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
log.debug("push: new work for nudge target {} raced the reminder loop's stop — restarting", lead);
scheduleNext(lead);
return;
}
log.debug("push: reminder loop ended for lead {}", lead);
log.debug("push: reminder loop ended for nudge target {}", lead);
}
/** Send one combined nudge covering everything currently pending for {@code lead}. */
@@ -773,13 +788,13 @@ public final class ReplyPushLoop {
List<PendingUnmappedTarget> unmappedTargets = pendingUnmappedTargetsFor(lead);
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty() && incidents.isEmpty()
&& unmappedTargets.isEmpty()) {
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
log.debug("push: pending work for nudge target {} drained before the nudge could be sent", lead);
return;
}
String nudge = formatNudge(replyTargets, tickets, questions, incidents, unmappedTargets);
try {
agents.send(lead, nudge);
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
log.debug("push: nudge sent to nudge target {} (reply {}/{}, ticket {}/{}, question {}/{}; "
+ "{} reply target(s), {} ticket(s), {} question(s))",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders,
@@ -796,7 +811,7 @@ public final class ReplyPushLoop {
}
}
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
log.warn("push: failed to nudge target {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders, e.toString());
}
@@ -814,7 +829,8 @@ public final class ReplyPushLoop {
List<PendingQuestion> questions, List<PendingIncident> incidents,
List<PendingUnmappedTarget> unmappedTargets) {
for (String target : replyTargets) {
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
pendingReplies.computeIfPresent(target,
(t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1, e.mayDrainNudge()));
}
for (PendingTicket ticket : tickets) {
pendingTickets.computeIfPresent(ticket.ticket(),
@@ -823,7 +839,7 @@ public final class ReplyPushLoop {
for (PendingQuestion question : questions) {
pendingQuestions.computeIfPresent(question.turnId(), (id, e) ->
new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(),
e.nudgeCount() + 1));
e.nudgeCount() + 1, e.mayAnswerNudge()));
}
for (PendingIncident incident : incidents) {
pendingIncidents.computeIfPresent(incident.key(), (id, e) -> new PendingIncident(e.key(),
@@ -285,13 +285,26 @@ public final class FleetApp {
/**
* 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.
* {@code observerSendTarget()} 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);
Predicate<String> observerSendTarget) {
return Authz.permits(caller, action, target, knownLeadOrCollaborator, observerSendTarget);
}
/**
* As {@link #permitsFor(Principal, Authz.Action, String, Predicate, Predicate)}, also
* threading the classifier an observer's {@code TASK_READ} is checked against; pass a real
* {@code MessageService#ownsTicket(String, String)}-backed predicate to exercise the real
* production gate for a ticket poll, as {@link #allow} does.
*/
static boolean permitsFor(Principal caller, Authz.Action action, String target,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> observerSendTarget,
Predicate<String> observerOwnsTicket) {
return Authz.permits(caller, action, target, knownLeadOrCollaborator, observerSendTarget,
observerOwnsTicket);
}
/**
@@ -303,11 +316,22 @@ public final class FleetApp {
* yours" (a worker reaching for another worker's session, or for orchestration).
*/
private boolean allow(Context ctx, Authz.Action action, String target) {
return allow(ctx, action, target, Authz.NO_OBSERVER_OWNED_TICKET);
}
/**
* As {@link #allow(Context, Authz.Action, String)}, with a real classifier for an observer's
* {@code TASK_READ} — the only call site that can supply one is a ticket poll, which knows the
* ticket id and the caller's owner key.
*/
private boolean allow(Context ctx, Authz.Action action, String target,
Predicate<String> observerOwnsTicket) {
if (auth == null) {
return true; // legacy: authorization not enforced
}
Principal caller = ctx.attribute(CALLER);
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.sendableObserverTarget())) {
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.observerSendTarget(),
observerOwnsTicket)) {
if (action != Authz.Action.READ && action != Authz.Action.METRICS
&& action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
@@ -929,11 +953,15 @@ public final class FleetApp {
/** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */
private void taskStatus(Context ctx) {
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
String ticket = ctx.pathParam("ticket");
Principal caller = ctx.attribute(CALLER);
String callerOwner = caller == null ? null : caller.ownerKey();
// The gate must see the ticket id, not a literal null -- an observer's grant is confined
// to a ticket its own caller created.
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), ticket, t -> messages.ownsTicket(t, callerOwner))) {
return;
}
Principal caller = ctx.attribute(CALLER);
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.ownerKey());
MessageService.TaskView v = messages.poll(ticket, callerOwner);
if (v == null) {
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
return;
@@ -47,6 +47,28 @@ class AuthzTest {
+ "relies on this gate refusing it first");
}
/**
* {@code fleet_inbox} collects the mail queued for the caller's own pane, so it is gated on
* terminal ownership and on nothing else: every role may do it for itself, and no role may do
* it for another pane.
*/
@Test
void collectingAnInboxIsOnlyEverForTheCallersOwnPane() {
for (Principal self : new Principal[]{WORKER_A, ARCH_DESIGN, COLLABORATOR, OBSERVER}) {
assertTrue(Authz.permits(self, INBOX, self.terminal()),
self.role() + " must be able to collect the mail for its own pane");
assertFalse(Authz.permits(self, INBOX, "term_someone_else"),
self.role() + " must not be able to collect another pane's mail");
}
// A lead carries a pane too, so it collects its own mail on the same rule.
assertTrue(Authz.permits(Principal.leader("opus", "term_lead", 800), INBOX, "term_lead"));
// The unnamed primary owns no pane, so there is no inbox it could be asking for.
assertFalse(Authz.permits(PRIMARY, INBOX, "term_a"),
"a caller with no pane of its own has no inbox to collect");
assertFalse(Authz.permits(WORKER_A, INBOX, null),
"a missing terminal must never match an owner");
}
@Test
void orchestrationBelongsToThePrimaryAlone() {
for (Authz.Action a : new Authz.Action[]{SPAWN, STOP, SEND, DRAIN}) {
@@ -288,16 +310,17 @@ class AuthzTest {
}
/**
* 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.
* Every action beyond READ/METRICS/REPLY/ASK/INBOX/SEND, asserted denied for an observer.
* The exempt set is the three only-as-itself actions plus the two open reads. {@code SEND}
* and {@code TASK_READ} are excluded here and given their own matrices below, since — unlike
* every action in this loop — each one's grant is conditional on more than the caller's role
* alone.
*/
@Test
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskAndSend() {
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxSendAndTaskRead() {
for (Authz.Action a : Authz.Action.values()) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == SEND) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND
|| a == TASK_READ) {
continue;
}
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
@@ -305,6 +328,49 @@ class AuthzTest {
}
}
/**
* {@code TASK_READ} is conditional for an observer, not fixed: it may poll a ticket its own
* {@code fleet_send(wait:false)} created, confined by {@code observerOwnsTicket}, and nothing
* else — a session's status is not ticket-scoped, so no call site can ever supply a predicate
* that grants it, and the 3-/4-/5-argument convenience forms stay fail-closed for it too.
*/
@Test
void anObserverMayTaskReadOnlyATicketTheClassifierSaysItOwns() {
assertTrue(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
"the classifier accepting the ticket as the caller's own must grant TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "someone_elses_ticket",
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, Authz.NO_OBSERVER_SEND_TARGET, target -> false),
"the classifier refusing the ticket must deny TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket"),
"the three-argument convenience form fails closed, so TASK_READ is refused without "
+ "a real classifier");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR),
"the four-argument form still fails closed for TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET),
"the five-argument form still fails closed for TASK_READ, since it supplies no "
+ "observerOwnsTicket classifier either");
}
/**
* Control for the test above: every other action's result for an observer does not move when
* the ticket-ownership classifier does. Only {@code TASK_READ} is wired to it.
*/
@Test
void theTicketOwnershipClassifierMovesOnlyTaskReadForAnObserver() {
for (Authz.Action a : Authz.Action.values()) {
if (a == TASK_READ) {
continue;
}
assertEquals(
Authz.permits(OBSERVER, a, "term_observer"),
Authz.permits(OBSERVER, a, "term_observer", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
a + " must not depend on the ticket-ownership classifier at all");
}
}
@Test
void anObserverIsNotCountedAsAnyOtherRole() {
assertFalse(OBSERVER.isPrimary());
@@ -912,42 +912,65 @@ class CallerResolverTest {
assertFalse(r.knownLeadOrCollaborator().test("term_a"));
}
// ── fleetd #743: sendableObserverTarget() reads the same maps and functions resolve() does ────
// ── observerSendTarget() reads the same maps and functions resolve() does, in the same order ──
/**
* 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() {
void observerSendTargetIsTrueForATerminalKnownAsNoOtherRole() {
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"));
assertTrue(r.observerSendTarget().test("term_other"));
}
/**
* A configured lead's own terminal is reachable: an observer may open a conversation with a
* lead, and this is the classifier that grant is checked against.
*/
@Test
void sendableObserverTargetIsFalseForALeadTerminal() {
void observerSendTargetIsTrueForALeadTerminal() {
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");
assertTrue(r.observerSendTarget().test("term_lead"),
"a lead's own terminal must be reachable from an observer pane");
}
/**
* A pane named as a lead AND bound to an architect slot resolves as the lead, because
* {@code resolve} reads the lead map first — so the classifier must accept it, or a pane's
* resolved role and its reachability would disagree.
*/
@Test
void observerSendTargetIsTrueForALeadTerminalThatIsAlsoABoundArchitectSlot() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_a", "opus-5.0"),
boundMembers("architect:lead-designer", MemberRole.ARCHITECT), t -> null, Map::of);
assertEquals(Role.PRIMARY, r.resolve("127.0.0.1", 42, null).role(),
"premise: the lead map is read before the architect registry");
assertTrue(r.observerSendTarget().test("term_a"));
}
@Test
void sendableObserverTargetIsFalseForACollaboratorTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
void observerSendTargetIsFalseForACollaboratorTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_lead", "opus-5.0"),
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");
assertFalse(r.observerSendTarget().test("term_collab"),
"a collaborator's own terminal must never be reachable from an observer pane");
// CONTROL: the same wiring, a target recognised as no configured role at all.
assertTrue(r.observerSendTarget().test("term_other"));
}
@Test
void sendableObserverTargetIsFalseForALiveSpawnedMembersTerminal() {
void observerSendTargetIsFalseForALiveSpawnedMembersTerminal() {
// 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,
@@ -958,26 +981,47 @@ class CallerResolverTest {
default -> null;
}, Map::of);
assertFalse(r.sendableObserverTarget().test("term_worker"));
assertFalse(r.sendableObserverTarget().test("term_architect"));
assertFalse(r.observerSendTarget().test("term_worker"));
assertFalse(r.observerSendTarget().test("term_architect"));
// CONTROL: the same wiring, a target the member lookup above answers null for.
assertTrue(r.observerSendTarget().test("term_other"));
}
/**
* A lead map entry does not rescue a terminal a live spawned member occupies: the member
* lookup runs first, exactly as in {@code resolve}.
*/
@Test
void observerSendTargetIsFalseForASpawnedMemberOnATerminalTheLeadMapAlsoNames() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_worker", "opus-5.0"), new MemberRegistry(null),
t -> "term_worker".equals(t) ? MemberRole.DEV : null, Map::of);
assertFalse(r.observerSendTarget().test("term_worker"));
// CONTROL: the same wiring, the same lead map, a terminal no member occupies.
assertTrue(r.observerSendTarget().test("term_other"));
}
@Test
void sendableObserverTargetIsFalseForABoundArchitectSlotWithNoLiveMember() {
void observerSendTargetIsFalseForABoundArchitectSlotWithNoLiveMember() {
// 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"));
assertFalse(r.observerSendTarget().test("term_a"));
// CONTROL: the same wiring, a terminal the bind above never touched.
assertTrue(r.observerSendTarget().test("term_other"));
}
@Test
void sendableObserverTargetIsFalseForANullTarget() {
void observerSendTargetIsFalseForANullTarget() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, Map::of);
assertFalse(r.sendableObserverTarget().test(null));
assertFalse(r.observerSendTarget().test(null));
// CONTROL: the same wiring, a non-null target.
assertTrue(r.observerSendTarget().test("term_other"));
}
}
@@ -254,10 +254,11 @@ public final class FakeHerdr implements HerdrClient {
}
/**
* 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.
* Override the text {@code agent.read} returns for the {@code detection} and {@code visible}
* sources only — the prompt/footer tail herdr uses for status detection and input-box probing, 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;
@@ -413,8 +414,8 @@ public final class FakeHerdr implements HerdrClient {
}
case "agent.read" -> {
Object source = params instanceof Map<?, ?> m ? m.get("source") : null;
String text = "detection".equals(source) && detectionText != null
? detectionText : readText;
boolean probeSource = "detection".equals(source) || "visible".equals(source);
String text = probeSource && detectionText != null ? detectionText : readText;
yield mapper.readTree(mapper.writeValueAsString(
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", text))));
}
@@ -12,6 +12,18 @@ class PromptBoxTest {
private static final String EMPTY = FakeHerdr.IDLE_PROMPT_CARET;
private static final String DRAFTED = FakeHerdr.DRAFTED_PROMPT_CARET;
/** An empty box: the pane draws a placeholder hint (its last submitted prompt) at the caret, faint. */
private static final String HINT_1 = "❯\u00a0\u001b[0m\u001b[2mstart md2trilium work for #2\u001b[0m";
/** Same shape as {@link #HINT_1}, a different placeholder hint. */
private static final String HINT_2 = "❯\u00a0\u001b[0m\u001b[2myes, push them\u001b[0m";
/** An empty box with no hint: the caret itself is drawn grey, with escape codes before the marker. */
private static final String EMPTY_GREY_CARET = "\u001b[0m\u001b[38;2;153;153;153m❯\u00a0\u001b[0m";
/** An empty box with no styling at all. */
private static final String EMPTY_PLAIN = "❯\u00a0";
// --- pure classification -------------------------------------------------
@Test
@@ -57,7 +69,7 @@ class PromptBoxTest {
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");
"the probed 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());
}
@@ -90,18 +102,54 @@ class PromptBoxTest {
"that marker survives in scrollback, and holding on it would hold every delivery forever");
}
@Test
void aPlaceholderHintReadsAsAnEmptyBox() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(HINT_1),
"the hint is the pane's own last prompt, drawn faint — it is not the operator's typing");
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(HINT_2));
}
@Test
void anEmptyBoxWithAGreyCaretReadsAsEmpty() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(EMPTY_GREY_CARET),
"the caret's own colour sits before the marker and must not stop the marker matching");
}
@Test
void anEmptyBoxWithNoStylingAtAllReadsAsEmpty() {
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(EMPTY_PLAIN));
}
@Test
void typedTextWithNoStylingIsADraftAndCountsItsCharacters() {
PromptBox.Reading reading = PromptBox.classify("❯\u00a0deploy the thing");
assertEquals(PromptBox.State.DRAFT, reading.state());
assertEquals("deploythething".length(), reading.characters());
}
@Test
void typedTextAfterAFaintHintCountsOnlyTheTextOutsideTheFaintSpan() {
PromptBox.Reading reading =
PromptBox.classify("❯\u00a0\u001b[0m\u001b[2mhint\u001b[0m and typed");
assertEquals(PromptBox.State.DRAFT, reading.state());
assertEquals("andtyped".length(), reading.characters(),
"the faint hint is excluded; only \"and typed\" was drawn plain");
}
// --- the gate ------------------------------------------------------------
@Test
void anEmptyBoxClearsTheGateAndReadsTheDetectionRegion() {
void anEmptyBoxClearsTheGateAndReadsTheVisibleRegionWithStylingKept() {
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");
assertEquals("visible", params.get("source"),
"the input box is drawn in the visible region, not in transcript scrollback");
assertEquals(false, params.get("strip_ansi"),
"styling must survive the read, or a faint placeholder hint reads as plain typed text");
}
@Test
@@ -0,0 +1,200 @@
package dev.ltms.fleet.inject;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.msg.TestTurnTokens;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Which route a message takes: offered to a pane that collects its own mail, or typed into the
* pane's terminal. Driven by feeding {@code onStatus}, so no real polling is involved.
*/
class InjectorModServedDeliveryTest {
/** A pane that collects its own mail. */
private static final String MOD = "term_mod";
/** A pane that does not, used as the control for every "nothing was typed" assertion. */
private static final String PTY = "term_pty";
private final FakeHerdr herdr = new FakeHerdr();
private final AtomicLong clock = new AtomicLong(1_000_000);
private final Injector injector = new Injector(new AgentControl(herdr), TurnListener.NOOP,
_ -> true, _ -> {
}, clock::get);
/** The messages typed into a pane, in order. A collected message must never appear here. */
@SuppressWarnings("unchecked")
private List<String> typed() {
return herdr.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.map(c -> ((Map<String, Object>) c.params()).get("text").toString())
.toList();
}
@Test
void aPaneThatCollectsItsOwnMailIsNeverTypedInto() {
injector.collectInbox(MOD); // the pane says it collects its own mail
CompletableFuture<Void> delivered =
injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD)).completion();
injector.onStatus(MOD, AgentStatus.IDLE); // offers it for collection
assertFalse(delivered.isDone(), "an offered message has not reached the pane yet");
assertEquals(List.of("do the task"), injector.collectInbox(MOD), "the pane collects it");
injector.onStatus(MOD, AgentStatus.IDLE); // the next sample records the delivery
assertTrue(delivered.isDone(), "a collected message is a delivered message");
assertEquals(List.of(), typed(), "nothing was typed into a pane that collects its own mail");
// The control: without it, an injector that typed nothing anywhere would pass the line
// above. Same injector, same herdr, a pane that never collected its mail.
injector.enqueue(PTY, "type this", TestTurnTokens.inert(PTY));
injector.onStatus(PTY, AgentStatus.IDLE);
assertEquals(List.of("type this"), typed(), "control: an ordinary pane is still typed into");
}
@Test
void aPaneThatStopsCollectingHasItsMailTypedInstead() {
injector.collectInbox(MOD);
CompletableFuture<Void> delivered =
injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD)).completion();
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of(), typed(), "control: while it is still collecting, nothing is typed");
assertFalse(delivered.isDone(), "control: and nothing is reported delivered either");
// The mod stopped calling fleet_inbox, so the pane leaves the window.
clock.addAndGet(PaneInbox.MOD_SERVED_WINDOW_MILLIS + 1);
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("do the task"), typed(), "the message falls back to the terminal route");
assertTrue(delivered.isDone(), "and is reported delivered once it is typed");
assertEquals(List.of(), injector.collectInbox(MOD),
"a message that was typed must not also still be collectable");
}
@Test
void aMessageOfferedForCollectionIsStillReportedAsNotYetDelivered() {
injector.collectInbox(MOD);
injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertFalse(injector.queuedWaitMillis(MOD) == null,
"an offered-but-uncollected message is still waiting, not delivered");
injector.collectInbox(MOD);
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(null, injector.queuedWaitMillis(MOD),
"once collected it is off the queue, the same as a typed message");
}
@Test
void aCollectedMessageIsNotFollowedByAnEnterNudge() {
injector.collectInbox(MOD);
injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
injector.collectInbox(MOD);
injector.onStatus(MOD, AgentStatus.IDLE); // records the delivery, arms the pickup latch
injector.onStatus(MOD, AgentStatus.IDLE); // still idle: the typed route nudges Enter here
assertFalse(herdr.called("agent.send_keys"),
"a pane that collects its own mail submits it itself; an Enter there would submit "
+ "whatever its operator is typing");
// The control: the nudge really does fire on the typed route, so the absence above is
// this route's behaviour and not a harness that never nudges at all.
injector.enqueue(PTY, "type this", TestTurnTokens.inert(PTY));
injector.onStatus(PTY, AgentStatus.IDLE);
injector.onStatus(PTY, AgentStatus.IDLE);
assertTrue(herdr.called("agent.send_keys"), "control: a typed message is nudged");
}
@Test
void aCancelledMessageStopsBeingCollectable() {
injector.collectInbox(MOD);
Injector.Delivery delivery = injector.enqueue(MOD, "retracted", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE); // offered for collection
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(delivery),
"an offered message has not reached the pane, so it can still be cancelled");
assertEquals(List.of(), injector.collectInbox(MOD),
"a cancelled message the caller was told never arrived must not arrive later");
// The control: an uncancelled message on the same route really is collectable, so the
// empty list above is the cancel working and not the offer never being made.
injector.enqueue(MOD, "kept", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("kept"), injector.collectInbox(MOD), "control: an offer is collectable");
}
@Test
void aMessageThePaneAlreadyCollectedCannotBeCancelled() {
injector.collectInbox(MOD);
// The control first: an offer the pane has not taken really is cancellable, so the
// different answer below is the collection and not a cancel that gave up on this route.
Injector.Delivery untaken = injector.enqueue(MOD, "retracted", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(untaken),
"control: an uncollected offer is still cancellable");
Injector.Delivery taken = injector.enqueue(MOD, "do the task", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("do the task"), injector.collectInbox(MOD), "the pane takes the offer");
// No poll has run since the pane took it, so the entry is still at the head and still
// QUEUED: the state alone cannot tell this case from an uncollected offer.
assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(taken),
"the pane holds this text and will act on it, so nothing can be cancelled");
injector.onStatus(MOD, AgentStatus.IDLE);
assertTrue(taken.completion().isDone(), "the next poll records the delivery");
assertEquals(List.of(), typed(), "and nothing was typed into the pane");
}
@Test
void cancellingALaterMessageLeavesACollectedOneDeliveredOnce() {
injector.collectInbox(MOD);
Injector.Delivery first = injector.enqueue(MOD, "first", TestTurnTokens.inert(MOD));
Injector.Delivery second = injector.enqueue(MOD, "second", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE); // offers the head
assertEquals(List.of("first"), injector.collectInbox(MOD), "the pane takes the head");
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(second),
"a message behind the collected one was never offered, so it cancels");
injector.onStatus(MOD, AgentStatus.IDLE);
assertTrue(first.completion().isDone(),
"cancelling a later message must not lose the record that the head was taken");
assertEquals(List.of(), injector.collectInbox(MOD),
"and the head must not be offered a second time");
// The control: the same injector still hands a later message over, so the empty
// collection above is this one not being re-offered rather than the route going quiet.
injector.onStatus(MOD, AgentStatus.WORKING); // the pane picks the collected message up
injector.onStatus(MOD, AgentStatus.IDLE); // and that turn ends
injector.enqueue(MOD, "third", TestTurnTokens.inert(MOD));
injector.onStatus(MOD, AgentStatus.IDLE);
assertEquals(List.of("third"), injector.collectInbox(MOD), "control: a later message is offered");
assertEquals(List.of(), typed(), "nothing took the terminal route");
}
@Test
void aPaneThatNeverCollectedIsTypedIntoFromTheStart() {
injector.enqueue(PTY, "do the task", TestTurnTokens.inert(PTY));
injector.onStatus(PTY, AgentStatus.IDLE);
assertEquals(List.of("do the task"), typed(), "no poll, no offer: the terminal route applies");
assertFalse(injector.isModServed(PTY));
}
}
@@ -0,0 +1,93 @@
package dev.ltms.fleet.inject;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/** What a pane that collects its own mail may and may not see. */
class PaneInboxTest {
private static final String A = "term_a";
private static final String B = "term_b";
private final AtomicLong clock = new AtomicLong(1_000_000);
private final PaneInbox inbox = new PaneInbox(clock::get);
@Test
void aPaneCollectsItsOwnMailAndLeavesTheNextPanesWhereItIs() {
inbox.offer(A, "for a");
inbox.offer(B, "for b");
assertEquals(List.of("for a"), inbox.drain(A), "a pane sees its own message");
// The control for the assertion below: B's message really is there to be missed, so a
// drain that returned everything would have shown it above.
assertEquals(List.of("for b"), inbox.drain(B), "the other pane's message stayed put");
}
@Test
void aCollectedMessageIsNotHandedOverASecondTime() {
inbox.offer(A, "deliver once");
assertEquals(List.of("deliver once"), inbox.drain(A), "control: the first call hands it over");
assertEquals(List.of(), inbox.drain(A), "a collected message is gone from the inbox");
}
@Test
void messagesComeBackInTheOrderTheyWereOffered() {
inbox.offer(A, "first");
inbox.offer(A, "second");
assertEquals(List.of("first", "second"), inbox.drain(A));
}
@Test
void aPaneIsModServedOnlyWhileItKeepsCollecting() {
assertFalse(inbox.isModServed(A), "a pane that has never collected is not mod-served");
inbox.drain(A);
assertTrue(inbox.isModServed(A), "control: collecting is what makes a pane mod-served");
clock.addAndGet(PaneInbox.MOD_SERVED_WINDOW_MILLIS);
assertTrue(inbox.isModServed(A), "control: the window edge still counts as collecting");
clock.addAndGet(1);
assertFalse(inbox.isModServed(A), "a pane that stopped collecting leaves the window");
}
@Test
void withdrawingTakesBackOnlyWhatThePaneHasNotCollected() {
PaneInbox.Entry collected = inbox.offer(A, "already taken");
inbox.drain(A);
PaneInbox.Entry pending = inbox.offer(A, "not taken yet");
inbox.withdrawAll(A);
assertTrue(collected.taken(), "withdrawing must not un-deliver a collected message");
assertFalse(pending.taken(), "control: the uncollected entry was never handed over");
assertEquals(List.of(), inbox.drain(A), "a withdrawn message is no longer collectable");
}
@Test
void forgettingAPaneDropsBothItsMailAndItsPollRecord() {
inbox.drain(A);
inbox.offer(A, "for a");
assertTrue(inbox.isModServed(A), "control: the pane is mod-served and holding mail");
inbox.forget(A);
assertFalse(inbox.isModServed(A), "a gone pane must not look mod-served to the next one");
assertEquals(List.of(), inbox.drain(A), "a gone pane's mail does not outlive it");
}
@Test
void aMissingTerminalCollectsNothingAndIsNeverModServed() {
assertEquals(List.of(), inbox.drain(null));
assertEquals(List.of(), inbox.drain(" "));
assertFalse(inbox.isModServed(null));
}
}
@@ -56,6 +56,7 @@ class FleetMcpAuthzTest {
private final AgentControl agents = new AgentControl(herdr);
private Metrics metrics;
private FleetMcp mcp;
private MessageService messages;
@AfterEach
void close() {
@@ -82,7 +83,7 @@ class FleetMcpAuthzTest {
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> "tok");
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(),
messages = new MessageService(agents, new Injector(agents), new Rendezvous(),
new InMemoryReplyInbox());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
metrics = FleetMetrics.create(sessions, new InMemoryReplyInbox());
@@ -270,18 +271,19 @@ 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 -------------------------------
// --- 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.
* wiring denies an observer's {@code SEND} to a collaborator and to a member while granting
* it to a lead and to another unclassified pane, proving the refusals are the rule and not a
* missing fixture.
*/
@Test
void anObserverMaySendOnlyToAnotherObserverNeverToALeadWorkerOrArchitect() {
void anObserverMaySendToALeadOrAnotherObserverButNeverToACollaboratorOrAMember() {
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),
@@ -289,22 +291,23 @@ class FleetMcpAuthzTest {
() -> 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");
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_lead_known"),
"an observer must reach a lead's terminal, so a peer session can open a "
+ "conversation with a lead");
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.
// configured roles above -- this is what proves the two 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
* As {@link #anObserverMaySendToALeadOrAnotherObserverButNeverToACollaboratorOrAMember}, 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.
*/
@@ -324,6 +327,57 @@ class FleetMcpAuthzTest {
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_other_observer"));
}
/**
* An observer's reach is widened for {@code SEND} alone. It still holds no {@code TASK_READ},
* so it cannot poll a ticket or read a lead's session status, and no {@code SPAWN}/
* {@code STOP}/{@code COORD_SEND}, so it cannot drive the fleet it can now message.
*/
@Test
void anObserverReachingALeadStillHoldsNoTicketReadAndNoLifecycle() {
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 -> null, Map::of);
FleetMcp m = mcp(true, callers);
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_lead_known"),
"premise: this wiring grants the observer's send to that lead");
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, "term_lead_known"));
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SPAWN, null));
assertNotNull(m.denyFor(OBSERVER, Authz.Action.STOP, null));
assertNotNull(m.denyFor(OBSERVER, Authz.Action.COORD_SEND, null));
}
/**
* An observer may read a ticket its own {@code fleet_send(wait:false)} created, never one
* another caller created — exercised against the real {@link MessageService#ownsTicket}, not
* a stand-in classifier.
*/
@Test
void anObserverMayTaskReadATicketItCreatedButNotOneAnotherCallerCreated() {
FleetMcp m = mcp(true);
String ownTicket = messages.sendAsync("term_a", "do it", null, OBSERVER);
String othersTicket = messages.sendAsync("term_a", "do it too", null, WORKER_A);
assertNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, ownTicket,
t -> messages.ownsTicket(t, OBSERVER.ownerKey())),
"an observer must read a ticket its own send created");
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, othersTicket,
t -> messages.ownsTicket(t, OBSERVER.ownerKey())),
"an observer must not read a ticket a different caller created");
}
/**
* {@code fleet_status} stays refused for an observer even once {@code TASK_READ} opens: it is
* not ticket-scoped, so no call site supplies a real {@code observerOwnsTicket} classifier,
* and the default 3-argument {@code denyFor} fails closed.
*/
@Test
void anObserverMayNotReadSessionStatusEvenAfterTaskReadOpensForTickets() {
FleetMcp m = mcp(true);
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, "term_a"),
"a session id is not a ticket the caller could ever own");
}
@Test
void theLegacyConstructorLeavesTheGateOpen() {
// The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced.
@@ -477,9 +531,10 @@ class FleetMcpAuthzTest {
/**
* {@link FleetMcp#leadsVisibleTo} is the whole policy decision for {@code fleet_list}'s
* {@code leads} array: visible to exactly the roles that may {@code SEND} to a lead -- the
* primary, an architect, and a collaborator -- never a worker, which holds {@code READ} but
* can never {@code SEND} at all, and never an anonymous caller.
* {@code leads} array: visible to the primary, an architect, and a collaborator -- never a
* worker, which holds {@code READ} but can never {@code SEND} at all, never an observer, which
* reads a lead's address from its filtered {@code panes} rows instead, and never an anonymous
* caller.
*/
@Test
void primaryArchitectAndCollaboratorMaySeeTheLeadsArray() {
@@ -489,6 +544,9 @@ class FleetMcpAuthzTest {
"a collaborator may SEND to a lead, so it must see the leads array to learn where");
assertFalse(FleetMcp.leadsVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the leads array");
assertFalse(FleetMcp.leadsVisibleTo(Principal.observer("term_obs", 700)),
"an observer may SEND to a lead but learns the address from its panes rows, which "
+ "carry no lead name, context window or config dir");
assertFalse(FleetMcp.leadsVisibleTo(ANON), "authenticated as nothing must not see it either");
}
@@ -513,8 +571,8 @@ class FleetMcpAuthzTest {
/**
* {@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
* primary, an architect, a collaborator, and an observer (to a lead or another observer pane,
* with its rows filtered and reduced -- see {@code listFleet}) -- never a worker, never an
* anonymous caller.
*/
@Test
@@ -525,8 +583,8 @@ class FleetMcpAuthzTest {
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");
"an observer holds SEND to a lead and 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");
}
@@ -0,0 +1,178 @@
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.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.inject.MemberPresence;
import dev.ltms.fleet.inject.TurnListener;
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.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Predicate;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* An observer's {@code fleet_send} to a lead, 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}, through the real {@link MessageService} and the real
* {@link Injector} to the point of its real herdr call.
*
* <p>{@link FleetMcpAuthzTest} proves {@code denyFor} grants this case. This is the path that also
* proves the grant is not dead at the injector's readiness gate: that gate is the production
* {@link Fleetd#deliverableTo} predicate, and the lead's terminal carries no
* {@link MemberPresence} entry, so the delivery can only pass by the lead being a configured lead.
*/
class FleetMcpObserverSendToLeadDeliveryTest {
private static final String LEAD = "term_lead_pane";
private static final String COLLABORATOR = "term_collab_pane";
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
private final MemberPresence presence = new MemberPresence();
private Injector injector;
private MessageService messages;
private FleetMcp mcp;
private Server server;
private String baseUrl;
@BeforeEach
void startServer() 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 -- so the real resolver lands the caller on the observer floor.
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 9001L);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> Map.of(LEAD, "fleet01-lead"), new MemberRegistry(null),
_ -> null, () -> Map.of(COLLABORATOR, "ops"));
Predicate<String> deliverable =
Fleetd.deliverableTo(presence, callers::leads, callers::collaborators);
injector = new Injector(agents, TurnListener.NOOP, deliverable);
messages = new MessageService(agents, injector, rendezvous, new InMemoryReplyInbox());
mcp = new FleetMcp(messages, workers, sessions, identity, presence,
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();
baseUrl = "http://127.0.0.1:"
+ ((ServerConnector) server.getConnectors()[0]).getLocalPort();
}
@AfterEach
void tearDown() throws Exception {
if (server != null) {
server.stop();
}
if (mcp != null) {
mcp.close();
}
}
@Test
void anObserversSendToALeadIsAttributedAndReachesTheRealInjector() throws Exception {
assertFalse(presence.isPresent(LEAD),
"premise: the lead's terminal is deliverable only as a configured lead, never "
+ "through a presence entry");
McpSchema.CallToolResult result = sendFleetSend(LEAD, "can we split the review?");
assertFalse(result.isError(), "an observer sending to a lead must be accepted: "
+ textOf(result));
long waiterDeadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(LEAD) && System.currentTimeMillis() < waiterDeadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(LEAD), "the async send must have opened its rendezvous waiter");
injector.onStatus(LEAD, 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]\ncan we split the review?",
params.get("text"),
"a lead must see the sender's own daemon-resolved terminal, never a raw echo of "
+ "the content and never a client-supplied name");
}
/**
* Control for the test above, over the same live server and the same wiring: the grant is
* specific to a lead target, so a collaborator's terminal is still refused at the handler and
* nothing is ever queued for it.
*/
@Test
void thatSameObserverIsStillRefusedACollaboratorsTerminal() {
McpSchema.CallToolResult result = sendFleetSend(COLLABORATOR, "can we split the review?");
assertTrue(result.isError(), "an observer must not reach a collaborator's terminal");
assertFalse(rendezvous.isWaiting(COLLABORATOR),
"a refused send must never open a waiter for its target");
}
private McpSchema.CallToolResult sendFleetSend(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();
}
}
@@ -1239,7 +1239,7 @@ class FleetMcpTest {
/**
* 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}
* the same classifier {@link CallerResolver#observerSendTarget} refuses as a {@code SEND}
* target — rather than falling through to {@code "observer"}.
*/
@Test
@@ -1276,14 +1276,13 @@ class FleetMcpTest {
}
/**
* 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},
* An observer's {@code fleet_list} carries a {@code panes} key, filtered to
* {@link CallerResolver#observerSendTarget} (so a lead's pane survives, while 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}.
* {@code workspaceLabel} (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() {
@@ -1306,7 +1305,7 @@ class FleetMcpTest {
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
() -> new PaneLocator(h).tabLabelsByTabId(),
() -> new PaneLocator(h).workspaceLabelsByWorkspaceId(), _ -> true,
callers::boundToArchitectSlot, callers.sendableObserverTarget());
callers::boundToArchitectSlot, callers.observerSendTarget());
String out = textOf(listFleetAsObserver(h, callers.leads(),
Map.of("term_collab_pane", "ops"), panes));
@@ -1314,7 +1313,10 @@ class FleetMcpTest {
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);
assertTrue(out.contains("\"sessionId\":\"term_lead_pane\""),
"a lead's pane is where an observer reads the sessionId its send needs: " + out);
assertTrue(out.contains("\"role\":\"lead\""),
"the lead's row must name the role, so an observer can tell it from a peer pane: " + 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"),
@@ -266,4 +266,61 @@ class PrimaryRegistryTest {
assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_never_seen").orElseThrow());
}
// ── fleetd #778: a nudge must offer only what the delegator's own role may run ──────────────
@Test
void theShortDelegationOverloadsGrantEverythingAPrimaryMayRun() {
var reg = new PrimaryRegistry(null);
reg.recordDelegation("term_worker", "term_lead");
assertEquals(Boolean.TRUE, reg.nudgeMayDrainFor("term_worker").orElseThrow());
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_worker").orElseThrow());
}
@Test
void theFullOverloadCarriesTheCallersOwnGrants() {
var reg = new PrimaryRegistry(null);
reg.recordDelegation("term_worker", "term_observer", "observer-name", false, false);
assertEquals(Boolean.FALSE, reg.nudgeMayDrainFor("term_worker").orElseThrow());
assertEquals(Boolean.FALSE, reg.nudgeMayAnswerFor("term_worker").orElseThrow());
}
@Test
void theTwoGrantsAreIndependent() {
var reg = new PrimaryRegistry(null);
reg.recordDelegation("term_worker", "term_architect", "architect-name", false, true);
assertEquals(Boolean.FALSE, reg.nudgeMayDrainFor("term_worker").orElseThrow(),
"an architect may not DRAIN");
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_worker").orElseThrow(),
"an architect may ANSWER");
}
@Test
void theSingletonPrimaryFallbackIsAlwaysGrantedEverything() {
var reg = new PrimaryRegistry("term_pinned");
assertEquals(Boolean.TRUE, reg.nudgeMayDrainFor("term_never_seen").orElseThrow());
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_never_seen").orElseThrow());
}
@Test
void withNoDelegationAndNoPrimaryTheGrantsAreUnknown() {
var reg = new PrimaryRegistry(null);
assertTrue(reg.nudgeMayDrainFor("term_never_seen").isEmpty());
assertTrue(reg.nudgeMayAnswerFor("term_never_seen").isEmpty());
}
@Test
void forgettingADelegationForgetsItsGrantsToo() {
var reg = new PrimaryRegistry(null);
reg.recordDelegation("term_worker", "term_observer", "observer-name", false, false);
reg.forgetDelegation("term_worker");
assertTrue(reg.nudgeMayDrainFor("term_worker").isEmpty());
assertTrue(reg.nudgeMayAnswerFor("term_worker").isEmpty());
}
}
@@ -0,0 +1,142 @@
package dev.ltms.fleet.msg;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.Injector;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* A delegation whose message the pane collected itself must reach the same ticket states as one
* typed into the pane.
*/
class MessageServiceInboxDeliveryTest {
/** A pane that collects its own mail. */
private static final String MOD = "term_mod";
/** A pane that does not — the control for every route-specific assertion here. */
private static final String PTY = "term_pty";
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
private final Injector injector = new Injector(agents);
private final InMemoryReplyInbox replyInbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox);
@BeforeEach
void setUp() {
replyInbox.own(MOD);
replyInbox.own(PTY);
}
/** The messages typed into a pane, in order. A collected message must never appear here. */
@SuppressWarnings("unchecked")
private List<String> typed() {
return herdr.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.map(c -> ((Map<String, Object>) c.params()).get("text").toString())
.toList();
}
/** {@code sendAsync} queues on another thread, so wait for the delivery to exist. */
private void awaitWaiting(String target) throws InterruptedException {
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(target),
"the delegation to " + target + " should have opened its rendezvous waiter");
}
/** A ticket's terminal state is stamped by the delegating thread, so poll until it settles. */
private MessageService.TaskView awaitSettled(String ticket) throws InterruptedException {
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 2000;
while ((view == null || view.phase() == MessageService.Phase.PENDING)
&& System.currentTimeMillis() < deadline) {
view = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(view, ticket + " must still be a known ticket");
return view;
}
@Test
void aTicketCollectedByThePaneReachesTheSameStatesAsATypedOne() throws Exception {
// The collecting pane announces itself before anything is delegated to it; the control
// pane never calls fleet_inbox at all.
assertEquals(List.of(), messages.collectInbox(MOD), "nothing is waiting yet");
String collectedTicket = messages.sendAsync(MOD, "long task");
awaitWaiting(MOD);
String typedTicket = messages.sendAsync(PTY, "long task");
awaitWaiting(PTY);
injector.onStatus(MOD, AgentStatus.IDLE); // offers the task for collection
injector.onStatus(PTY, AgentStatus.IDLE); // types the task into the pane
assertEquals(List.of("long task"), typed(),
"only the control pane was typed into; the collecting pane was not");
assertTrue(messages.poll(collectedTicket).detail().contains("queued, not yet delivered"),
"control: an offered-but-uncollected message has not reached its pane");
assertFalse(messages.poll(typedTicket).detail().contains("queued, not yet delivered"),
"control: a typed message has reached its pane");
assertEquals(List.of("long task"), messages.collectInbox(MOD), "the pane collects the task");
injector.onStatus(MOD, AgentStatus.IDLE); // records the delivery
assertFalse(messages.poll(collectedTicket).detail().contains("queued, not yet delivered"),
"a collected message is reported delivered, the same as a typed one");
injector.onStatus(MOD, AgentStatus.WORKING);
injector.onStatus(PTY, AgentStatus.WORKING);
// The reply still arrives through fleet_reply, and lands the same way on both routes.
MessageService.ReplyOutcome collectedReply = messages.reply(MOD, "task done");
MessageService.ReplyOutcome typedReply = messages.reply(PTY, "task done");
assertEquals(typedReply, collectedReply, "both routes resolve their delegation the same way");
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, collectedReply);
MessageService.TaskView collectedDone = awaitSettled(collectedTicket);
MessageService.TaskView typedDone = awaitSettled(typedTicket);
assertEquals(MessageService.Phase.DONE, collectedDone.phase(),
"a ticket delivered by collection must not strand at PENDING");
assertEquals(MessageService.Phase.DONE, typedDone.phase(),
"control: the typed route reaches the same terminal state");
assertEquals("task done", collectedDone.reply());
assertEquals(typedDone.replySource(), collectedDone.replySource());
assertEquals(List.of("long task"), typed(),
"the collected delegation ran start to finish without typing into its pane");
}
@Test
void collectingAnInboxReachesOnlyTheCallersOwnMail() throws Exception {
assertEquals(List.of(), messages.collectInbox(MOD));
assertEquals(List.of(), messages.collectInbox(PTY));
messages.sendAsync(MOD, "task for the mod pane");
awaitWaiting(MOD);
messages.sendAsync(PTY, "task for the other pane");
awaitWaiting(PTY);
injector.onStatus(MOD, AgentStatus.IDLE);
injector.onStatus(PTY, AgentStatus.IDLE);
assertEquals(List.of("task for the mod pane"), messages.collectInbox(MOD));
// The control: the other pane's task really was waiting for it, so a collect that
// returned everyone's mail would have shown it on the line above.
assertEquals(List.of("task for the other pane"), messages.collectInbox(PTY));
}
}
@@ -1034,6 +1034,27 @@ class MessageServiceTest {
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
}
/**
* {@link MessageService#ownsTicket(String, String)} is the side-effect-free ownership check
* the authorization gate consults ahead of {@link MessageService#poll(String, String)}, which
* also fires the ticket's collection hook. Read-only: polling the same ticket afterwards still
* sees it, which {@link #theOneArgPollOverloadBypassesOwnershipEntirely} nearby does not
* guarantee for every overload.
*/
@Test
void publicOwnsTicketIsASideEffectFreeOwnershipCheck() {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "long task", null, lead);
assertTrue(messages.ownsTicket(ticket, lead.ownerKey()));
assertFalse(messages.ownsTicket(ticket, Principal.anonymous().ownerKey()));
assertFalse(messages.ownsTicket("no-such-ticket", lead.ownerKey()),
"a ticket this daemon never heard of is owned by nobody");
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket, lead.ownerKey()).phase(),
"the ownership check must not have consumed or altered the ticket");
}
@Test
void architectOwnershipUsesTerminalRatherThanSlot() {
Principal oldArchitect = Principal.architect("opus", "term_OLD", 1);
@@ -339,6 +339,85 @@ class ReplyPushLoopTest {
assertEquals(0, rec.promptTargets().size(), "nobody live was found, so nothing was ever sent");
}
// --- fleetd #778: never nudge a delegator to run a call its own role is refused -------------
@Test
void onReplyQueuedNeverNudgesADelegatorThatMayNotRunDrain() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
loop(1, 50).onReplyQueued(WORKER);
Thread.sleep(200);
assertEquals(0, rec.sendCount(),
"a delegator whose role cannot run fleet_poll(target=...) must never be nudged to run it");
}
@Test
void onReplyQueuedStillNudgesADelegatorThatMayRunDrain() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "lead", true, true);
loop(1, 50).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"),
"a delegator that may run DRAIN must still be nudged about it: " + nudge);
}
@Test
void onQuestionOpenedNeverNudgesADelegatorThatMayNotRunAnswer() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
Thread.sleep(200);
assertEquals(0, rec.sendCount(),
"a delegator whose role cannot run fleet_send(turnId=...) must never be nudged to run it");
}
@Test
void onQuestionOpenedStillNudgesADelegatorThatMayRunAnswer() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
registry.recordDelegation(WORKER, OTHER_PRIMARY, "architect", false, true);
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("fleet_send(turnId="),
"an architect may run ANSWER, so the nudge must still name it: " + nudge);
}
/**
* A delegator barred from DRAIN can still have a ticket genuinely pending for it — tickets are
* out of this unit's scope (fleetd #778 Part 1 already makes a ticket nudge correct for an
* observer). The reply half must still be suppressed even though the ticket half fires.
*/
@Test
void aForbiddenReplyNudgeIsSuppressedWhileAnEligibleTicketStillFires() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
var loop = loop(1, 300); // wide backoff so both calls land before the first tick fires
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertFalse(nudge.contains("fleet_poll(target=" + WORKER + ")"),
"the reply half is forbidden for this delegator and must never be named: " + nudge);
}
// --- nudge format --------------------------------------------------------------------------
@Test
@@ -285,6 +285,67 @@ class FleetAppAuthTest {
}
}
/**
* An observer may read a ticket its own {@code fleet_send(wait:false)} created, over the REST
* route and not only the unit-level classifier — and still not a ticket a different caller
* created, even another observer pane's.
*/
@Test
void restPollLetsAnObserverReadItsOwnTicketButNotAnothersOverRest() 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 observerApp = startOnSharedServiceAsObserver(messages, herdr, 9001L); // -> term_shell
Javalin otherObserverApp = startOnSharedServiceAsObserver(messages, herdr, FakeHerdr.WORKER_PID); // -> term_a
try {
Principal observer = Principal.observer("term_shell", 9001);
String ticket = messages.sendAsync("term_a", "long task", null, observer);
HttpResponse<String> own = send(observerApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, own.statusCode());
assertFalse(own.body().contains("forbidden"),
"the observer that created the ticket must read it: " + own.body());
// Unlike a worker's or the primary's mismatched ticket (refused by MessageService's own
// ownership check, inside a 200 response), a different observer is refused at the
// authorization gate itself, since its grant is conditional on the classifier -- so
// this is a 403, never reaching poll().
HttpResponse<String> other = send(otherObserverApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(403, other.statusCode(), other.body());
assertFalse(other.body().contains("\"reply\""),
"a refusal must never carry reply text: " + other.body());
} finally {
observerApp.stop();
otherObserverApp.stop();
}
}
/**
* As {@link #startOnSharedService(MessageService, FakeHerdr, long)}, but resolves every
* connecting pane to the unconfigured-pane {@link Role#OBSERVER} floor instead of a spawned
* worker — no lead, collaborator, or member map claims it.
*/
private Javalin startOnSharedServiceAsObserver(MessageService messages, FakeHerdr herdr, long pid) {
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
AgentControl agents = new AgentControl(herdr);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
Map.of(wcfg.profile(), wcfg), wcfg.profile(),
k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null), t -> null, Map::of);
Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
callers, appMetrics).build().start("127.0.0.1", 0);
}
/**
* {@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
@@ -694,6 +755,27 @@ class FleetAppAuthTest {
assertEquals(403, toSpawnedMembersTerminal.statusCode(), toSpawnedMembersTerminal.body());
}
/**
* The REST route must give the same answer as MCP for an observer: pid 9001 resolves to
* "term_shell", a herdr pane recognised as no configured role, so the real
* {@link CallerResolver#observerSendTarget()} classifier reaches the known lead and refuses
* the known collaborator -- over the route, not just the unit-level classifier, so a grant
* covering only MCP cannot leave this one behind.
*/
@Test
void anObserverMaySendToAKnownLeadButNotToAKnownCollaboratorOverRest() throws Exception {
int port = startWithRealClassifier(9001L,
Map.of("term_lead_known", "lead-x"), Map.of("term_collab_known", "ops2"));
HttpResponse<String> toLead = send(port, "POST", "/sessions/term_lead_known/message",
"{\"content\":\"hi\",\"wait\":false}", null);
assertEquals(202, toLead.statusCode(), toLead.body());
HttpResponse<String> toCollaborator = send(port, "POST", "/sessions/term_collab_known/message",
"{\"content\":\"hi\",\"wait\":false}", null);
assertEquals(403, toCollaborator.statusCode(), toCollaborator.body());
}
/**
* As {@link #start}, but with explicit lead/collaborator maps and no spawned-member roster, so
* a test can wire the real {@link CallerResolver#knownLeadOrCollaborator()} classifier instead
+2 -2
View File
@@ -1,7 +1,7 @@
{
"name": "fleet",
"description": "Make a project fleet-ready: mount the fleetd MCP gateway and set up standard Claude Code settings so this session can orchestrate a fleet of delegated workers. Lead-side only — member skills and agents travel in the worktree. Ships no credentials.",
"version": "0.2.0",
"description": "Set up standard Claude Code settings so this session can orchestrate a fleet of delegated workers, and run the fleet mod for cross-session messaging. Lead-side only — member skills and agents travel in the worktree, and mounting the fleetd MCP gateway is now the instance's or the project's job, not this plugin's. Ships no credentials.",
"version": "0.3.0",
"author": {
"name": "LTMS"
},
-8
View File
@@ -1,8 +0,0 @@
{
"mcpServers": {
"fleet": {
"type": "http",
"url": "${FLEETD_MCP_URL}"
}
}
}
+35 -19
View File
@@ -1,7 +1,7 @@
# fleet (Claude Code plugin)
Makes a project **fleet-ready**: mounts the `fleetd` MCP gateway and applies standard Claude Code
settings, so the session can orchestrate a fleet of delegated workers.
Makes a project **fleet-ready**: applies standard Claude Code settings and runs the fleet mod, so
the session can orchestrate a fleet of delegated workers.
**This plugin ships no credentials.** Every secret is referenced by environment-variable *name*;
the values stay with the user. Nothing the plugin writes is unsafe to commit.
@@ -22,9 +22,10 @@ Member-facing assets travel in the worktree, not in this plugin. See fleetd #362
## What it is not
The plugin is the **client-side setup**, not the bridge. `fleetd` is a separate daemon and `herdr`
is a separate PTY multiplexer, each with its own lifecycle and install. The plugin mounts an
already-running daemon and tells you what is missing when one isn't there — it deliberately does
not try to install system services on your behalf.
is a separate PTY multiplexer, each with its own lifecycle and install, and the plugin does not try
to install either on your behalf. It also does not mount the daemon for you — mounting is the
instance's or the project's own `.mcp.json`, and `/fleet:setup` is the one thing in this plugin that
still helps with that (it writes the project-level entry).
## Install
@@ -33,11 +34,20 @@ not try to install system services on your behalf.
/plugin install fleet@fleetd
```
Export the gateway URL — the plugin mounts `${FLEETD_MCP_URL}`, not a hardcoded address, so one
plugin serves hosts that run the daemon on different ports:
Mount the daemon yourself — this plugin carries no mount of its own. Either add the entry below to
your Claude Code instance's own `.claude.json`, so every project you open there gets it, or run
`/fleet:setup` in the project you want to onboard, which writes the same entry into that project's
`.mcp.json`:
```shell
export FLEETD_MCP_URL=http://127.0.0.1:8765/mcp
```json
{
"mcpServers": {
"fleet": {
"type": "http",
"url": "http://127.0.0.1:8765/mcp"
}
}
}
```
Then, in the project you want to onboard:
@@ -50,18 +60,24 @@ Then, in the project you want to onboard:
| Component | Effect |
|---|---|
| `.mcp.json` | mounts `fleet` at `${FLEETD_MCP_URL}` for any session with the plugin enabled |
| `skills/setup` | `/fleet:setup` — preflight, project settings, credential guidance, and verification |
| `skills/setup` | `/fleet:setup` — preflight, project settings, credential guidance, and verification. Also the only thing in this plugin that helps mount `fleet`: it writes the project `.mcp.json` entry shown above. |
| `hooks/register.js` (the fleet mod) | Cross-account session messaging while the plugin is enabled: `/fleet-peers`, `/fleet-mail`, `/fleet-whoami`, and a background poll that delivers mail fleetd queued for this pane. A spawned worker or architect skips that poll, because it already gets its brief pasted into its pane. |
The server is named **`fleet`** on purpose: that is `PeerLauncher.MCP_MOUNT_NAME` in the daemon and
the name a spawned member's own mount carries. Version 0.1.0 named it `fleetd`, which produced two
mounts of one daemon for anyone who also had a project-level `.mcp.json`. Upgrading from 0.1.0 is a
**breaking change** — a project that pre-allowed `mcp__fleetd__fleet_whoami` in
`.claude/settings.json` must be updated to `mcp__fleet__*`.
Whichever file mounts the daemon, name the server **`fleet`**. That is `PeerLauncher.MCP_MOUNT_NAME`
in the daemon, the name a spawned member's own mount carries, and the name the `mcp__fleet__*`
role heuristic in `CLAUDE.md` keys on.
Because the plugin carries its own `.mcp.json`, an installed plugin needs no project-level MCP
file at all. The setup skill writes one only when you want the mount to work *without* the plugin —
for teammates who haven't installed it, or for CI.
## Upgrading from 0.2.0 — breaking
The plugin no longer mounts the daemon. It used to carry its own `.mcp.json`, pointed at
`${FLEETD_MCP_URL}`, and that file is gone. The mod still reads `FLEETD_MCP_URL`, but only as an
optional override of the address it calls (default `http://127.0.0.1:8765/mcp`). Mount `fleet`
yourself: add the entry under **Install** above to your instance's `.claude.json` or to the
project's own `.mcp.json`, by hand or with `/fleet:setup`.
Version 0.1.0 named the mounted server `fleetd`, which produced two mounts of one daemon for
anyone who also had a project-level `.mcp.json`. A project that pre-allowed
`mcp__fleetd__fleet_whoami` in `.claude/settings.json` must be updated to `mcp__fleet__*`.
## Verifying a setup
+4
View File
@@ -0,0 +1,4 @@
{
"description": "fleet mod: cross-account session messaging and PTY-free delivery",
"modules": ["./register.js"]
}
+281
View File
@@ -0,0 +1,281 @@
// fleet mod — session messaging for the claude-bridge fleet.
//
// $.store is kept under CLAUDE_CONFIG_DIR, so /fleet-peers and /fleet-mail reach
// only sessions that share this session's config dir. fleetd names a caller by
// its pane, not its account, so /fleet-whoami and anything else sent through
// fleetTool reach the fleet from either account.
//
// Delivery uses $.prompt.submit, so nothing is typed into a pane and no prompt
// box is read.
//
// There are two inboxes. The $.store one carries /fleet-mail between sessions
// that share this config dir. fleet_inbox carries what the fleet queued for
// this pane, from any account, and calling it is also what tells fleetd to
// queue here rather than type into the terminal.
const PRESENCE_PREFIX = 'presence:'
const INBOX_PREFIX = 'inbox:'
const PRESENCE_REFRESH_MS = 15_000
const INBOX_POLL_MS = 3_000
// fleetd stops queueing for a pane that goes quiet, so this must stay well
// under the window the daemon allows between calls.
const FLEETD_INBOX_POLL_MS = 3_000
// A session whose presence row is older than this is treated as gone. It must
// exceed PRESENCE_REFRESH_MS by enough that one missed refresh is not a death.
const PRESENCE_STALE_MS = 60_000
const MAX_INBOX = 50
/** The store key holding one session's queued messages. */
function inboxKey(sessionId) {
return INBOX_PREFIX + sessionId
}
/** The store key holding one session's presence row. */
function presenceKey(sessionId) {
return PRESENCE_PREFIX + sessionId
}
/**
* Append one message to a target's inbox.
*
* $.store has no compare-and-swap, so two senders writing in the same instant
* can lose a message. Callers that need delivery confirmed should read the
* inbox back.
*/
async function deliver($, target, message) {
const key = inboxKey(target)
const queued = (await $.store.get(key)) || []
queued.push(message)
// Keep the newest: an unread inbox must not grow without bound.
const kept = queued.slice(-MAX_INBOX)
await $.store.set(key, kept)
return kept.length
}
/** Every session that refreshed its presence row recently, newest first. */
async function livePeers($, now) {
const keys = await $.store.keys()
const rows = []
for (const key of keys) {
if (!key.startsWith(PRESENCE_PREFIX)) continue
const row = await $.store.get(key)
if (!row || typeof row.at !== 'number') continue
if (now - row.at > PRESENCE_STALE_MS) continue
rows.push(row)
}
rows.sort((a, b) => b.at - a.at)
return rows
}
const FLEETD_MCP_DEFAULT = 'http://127.0.0.1:8765/mcp'
const MCP_HEADERS = { 'Content-Type': 'application/json', Accept: 'application/json, text/event-stream' }
// The status fleetd answers, with "Session not found", for an Mcp-Session-Id it no longer holds.
const MCP_SESSION_GONE = 404
// The MCP session every call below shares, and the URL it was opened against. fleetd keeps a
// server-side session per initialize and drops it only on a DELETE, so one initialize per call
// would leave a session behind every time.
let mcpSessionId = null
let fleetdUrl = FLEETD_MCP_DEFAULT
/**
* Open an MCP session on the local fleetd and hold it for later calls.
*
* Cleared first, so a failure here leaves no dead id behind for the next call to reuse. Reads
* FLEETD_MCP_URL fresh on every open, so a session opened after the daemon moves uses the new
* address.
*/
async function openFleetSession($) {
mcpSessionId = null
fleetdUrl = (await $.env.get('FLEETD_MCP_URL')) || FLEETD_MCP_DEFAULT
const init = await $.http.fetch(fleetdUrl, {
method: 'POST',
headers: MCP_HEADERS,
body: JSON.stringify({
jsonrpc: '2.0', id: 1, method: 'initialize',
params: { protocolVersion: '2025-06-18', capabilities: {}, clientInfo: { name: 'fleet-mod', version: '0' } },
}),
})
if (!init.ok) throw new Error('fleetd initialize failed with status ' + init.status)
const opened = init.headers['mcp-session-id']
await $.http.fetch(fleetdUrl, {
method: 'POST',
headers: { ...MCP_HEADERS, 'Mcp-Session-Id': opened },
body: JSON.stringify({ jsonrpc: '2.0', method: 'notifications/initialized' }),
})
mcpSessionId = opened
}
/** Send one tools/call on the session this mod holds, and return the raw HTTP answer. */
function sendFleetToolCall($, tool, args) {
return $.http.fetch(fleetdUrl, {
method: 'POST',
headers: { ...MCP_HEADERS, 'Mcp-Session-Id': mcpSessionId },
body: JSON.stringify({ jsonrpc: '2.0', id: 2, method: 'tools/call', params: { name: tool, arguments: args } }),
})
}
/**
* Call one fleet_* tool on the local fleetd, and return its text result.
*
* fleetd names the caller from the TCP connection on every request, not from the MCP session, so
* the answer is about this session's own pane whichever Claude account the session runs on, and
* reusing one session never changes whose call it is.
*/
async function fleetTool($, tool, args) {
if (mcpSessionId === null) await openFleetSession($)
let call = await sendFleetToolCall($, tool, args)
if (call.status === MCP_SESSION_GONE) {
// A daemon restart drops every session it held. Open a new one and retry once.
await openFleetSession($)
call = await sendFleetToolCall($, tool, args)
}
if (!call.ok) throw new Error('fleetd ' + tool + ' failed with status ' + call.status)
// A tool call answers as a server-sent event: the JSON is on the data: line.
const line = call.text.split('\n').find((l) => l.startsWith('data:'))
const body = JSON.parse(line ? line.slice(5) : call.text)
if (body.error) throw new Error(body.error.message)
return body.result.content.map((c) => c.text).join('\n')
}
export function register(on) {
on('session.start', async ($, e, next) => {
const self = await $.session.id()
// Announce before the first refresh is due, or a session shorter than one
// refresh interval never appears to its peers at all.
await $.store.set(presenceKey(self), {
sessionId: self,
cwd: await $.session.cwd(),
at: await $.clock.now(),
})
// Keep announcing: a row that stops being refreshed is how another session
// learns this one is gone.
$.clock.every(PRESENCE_REFRESH_MS, async () => {
const now = await $.clock.now()
await $.store.set(presenceKey(self), {
sessionId: self,
cwd: await $.session.cwd(),
at: now,
})
})
// Collect this session's mail and hand it to Claude. $.prompt.submit waits
// for the session to be idle, so this never lands mid-turn.
$.clock.every(INBOX_POLL_MS, async () => {
const key = inboxKey(self)
const queued = (await $.store.get(key)) || []
if (queued.length === 0) return
await $.store.set(key, [])
for (const message of queued) {
await $.prompt.submit({
text: 'Message from fleet session ' + message.from + ':\n\n' + message.text,
})
}
})
// A spawned worker or architect already gets its brief pasted into its pane, so this
// poll would be a second, redundant delivery path for it. Every other role collects
// its own mail through this poll.
let role = null
// Collect what the fleet queued for this pane and hand each message to
// Claude. Every call also renews fleetd's record that this pane collects its
// own mail, so an empty answer still has to be asked for.
$.clock.every(FLEETD_INBOX_POLL_MS, async () => {
if (role === null) {
try {
role = JSON.parse(await fleetTool($, 'fleet_whoami', {})).role
} catch {
// A daemon that is down, or a pane fleetd cannot place, is the ordinary
// case on a host with no fleet running. Retry on the next tick.
return
}
}
if (role === 'worker' || role === 'architect') return
let collected
try {
collected = JSON.parse(await fleetTool($, 'fleet_inbox', {}))
} catch {
// The timer survives a throw, so this only keeps every tick from
// writing an error to the debug log.
return
}
for (const text of collected.messages || []) {
await $.prompt.submit({ text: 'Message from the fleet, via fleetd:\n\n' + text })
}
})
await $.command.register({
name: 'fleet-peers',
description: 'List fleet sessions on this machine, including other accounts',
})
await $.command.register({
name: 'fleet-whoami',
description: 'Show who fleetd says this session is',
})
await $.command.register({
name: 'fleet-mail',
description: 'Send a message to a fleet session on this machine',
argumentHint: '<sessionId> <text>',
// Runs even while Claude is working, so a correction is never queued
// behind the turn it is meant to correct.
immediate: true,
})
return next(e)
})
on('command.run', { command: 'fleet-peers' }, async ($) => {
const now = await $.clock.now()
const self = await $.session.id()
const peers = await livePeers($, now)
if (peers.length === 0) return { text: 'No fleet sessions have announced themselves yet.' }
const lines = peers.map((p) => {
const age = Math.round((now - p.at) / 1000)
const mark = p.sessionId === self ? ' (this session)' : ''
return p.sessionId + ' ' + p.cwd + ' seen ' + age + 's ago' + mark
})
return { text: 'Fleet sessions on this machine:\n' + lines.join('\n') }
})
on('command.run', { command: 'fleet-mail' }, async ($, e) => {
const args = (e.args || '').trim()
const split = args.indexOf(' ')
if (split < 1) return { text: 'Usage: /fleet-mail <sessionId> <text>' }
const target = args.slice(0, split)
const text = args.slice(split + 1).trim()
if (text === '') return { text: 'Usage: /fleet-mail <sessionId> <text>' }
const now = await $.clock.now()
const peers = await livePeers($, now)
if (!peers.some((p) => p.sessionId === target)) {
return { text: 'No live fleet session ' + target + '. Run /fleet-peers.' }
}
const self = await $.session.id()
const depth = await deliver($, target, { from: self, text: text, at: now })
return { text: 'Queued for ' + target + ' (' + depth + ' in its inbox).' }
})
on('command.run', { command: 'fleet-whoami' }, async ($) => {
try {
return { text: 'fleetd says: ' + (await fleetTool($, 'fleet_whoami', {})) }
} catch (err) {
return { text: 'fleetd unreachable: ' + err.message }
}
})
// Record what arrives over Claude Code's own channel, so a message delivered
// by the fleet and one delivered by SendMessage can be told apart.
//
// This hook gates delivery, so it must never decide the message's fate. The
// .catch handler passes the message on when the logging above throws.
on('session.receive', async ($, e, next) => {
$.ui.log('fleet: inbound ' + (e.origin && e.origin.kind) + ', ' + String(e.text).length + ' chars')
return next(e)
}).catch(async ($, e, next) => {
if (next.called) return undefined
return next(e)
})
}
+6 -12
View File
@@ -36,16 +36,11 @@ a time.
command -v herdr && herdr --version 2>&1 | head -1 || echo "MISSING: herdr"
command -v ccs && ccs version 2>&1 | head -1 || echo "MISSING: ccs (needed for worker profiles)"
command -v codex && codex --version 2>&1 | head -1 || echo "absent: codex (optional)"
curl -s -m 5 "${FLEETD_MCP_URL%/mcp}/healthz" 2>/dev/null \
|| curl -s -m 5 http://127.0.0.1:8765/healthz \
|| echo "MISSING: fleetd daemon is not reachable"
[ -n "$FLEETD_MCP_URL" ] && echo "FLEETD_MCP_URL is set" || echo "MISSING: FLEETD_MCP_URL"
curl -s -m 5 http://127.0.0.1:8765/healthz || echo "MISSING: fleetd daemon is not reachable"
```
**`FLEETD_MCP_URL` is required.** The plugin's own `.mcp.json` mounts `${FLEETD_MCP_URL}` rather
than a hardcoded address, so one plugin can serve hosts that run the daemon on different ports. If
it is unset the mount does not resolve. The usual value is `http://127.0.0.1:8765/mcp`; tell the
user to export it, do not write it into a file for them.
The usual address is `http://127.0.0.1:8765`. If this daemon runs elsewhere, use that address
instead wherever this skill writes `http://127.0.0.1:8765/mcp` below.
A healthy daemon answers with its status **and the herdr protocol it negotiated**:
@@ -95,10 +90,9 @@ If `.mcp.json` already exists, add only the `fleet` key and leave every other se
If a `fleet` entry is already there with a different URL, **ask** rather than assuming yours is
right — a non-default port usually means a deliberate second daemon.
> **If this plugin is installed, you can skip this step entirely.** The plugin ships its own
> `.mcp.json`, so `fleetd` is already mounted for any session with the plugin enabled. Write the
> project-level file only when the user wants the mount to work *without* the plugin — for
> teammates who have not installed it, or for CI.
This write is the only way this plugin helps mount `fleet` — the plugin carries no mount of its
own. A session that wants the mount without running this skill can instead add the same entry to
its own Claude Code instance's `.claude.json`.
**Before writing it, settle whether `.mcp.json` is committed here:**
+490
View File
@@ -0,0 +1,490 @@
import { expect, mock, test } from 'claude-code/testing'
// The store-backed tests write the row another session would write, because a
// test cannot start a second session.
const PEER = 'peer-session'
const SELF = 'this-session'
// A fixed clock keeps the staleness arithmetic exact.
const NOW = 1_700_000_000_000
/**
* Answer the mods API calls the harness has no implementation for, and stand in
* for Claude Code's own behaviour beneath the mod's gating hooks.
*
* Returns the map backing $.store. The harness puts no `store` namespace on the
* test's own `$`, so a test seeds and inspects the carrier through this map.
*/
function stubEngine(on: any): Map<string, any> {
const store = new Map<string, any>()
on('store.get', (_$: any, e: any) => ({ value: store.get(e.key) }))
on('store.set', (_$: any, e: any) => {
store.set(e.key, e.value)
return { value: undefined }
})
on('store.delete', (_$: any, e: any) => {
store.delete(e.key)
return { value: undefined }
})
on('store.keys', () => ({ value: [...store.keys()] }))
on('clock.now', () => ({ value: NOW }))
on('session.id', () => ({ value: SELF }))
on('session.cwd', () => ({ value: '/Users/x/claude-bridge' }))
on('ui.log', () => ({ value: undefined }))
on('prompt.submit', () => ({ value: undefined }))
on('env.get', () => ({ value: undefined }))
// Claude Code's own delivery, which the mod's receive hook must reach.
on('session.receive', (_$: any, e: any) => e)
return store
}
/** The presence row a session on the other account would write. */
function announce(store: Map<string, any>, sessionId: string, cwd: string, at: number) {
store.set('presence:' + sessionId, { sessionId, cwd, at })
}
test('/fleet-peers lists a session that announced itself', async ($, on) => {
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW)
const answer = await $.command.run({ command: 'fleet-peers', args: '' })
expect(answer.text).toContain(PEER)
expect(answer.text).toContain('/Users/x/work-repo')
})
test('/fleet-peers hides a session whose presence row went stale', async ($, on) => {
const store = stubEngine(on)
// One second past the 60s staleness cut.
announce(store, PEER, '/Users/x/work-repo', NOW - 61_000)
const answer = await $.command.run({ command: 'fleet-peers', args: '' })
expect(answer.text).not.toContain(PEER)
})
test('/fleet-peers keeps a row one second inside the staleness cut', async ($, on) => {
// The positive control for the test above: without this, a bug that hid
// every row would still satisfy that assertion.
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW - 59_000)
const answer = await $.command.run({ command: 'fleet-peers', args: '' })
expect(answer.text).toContain(PEER)
})
test('/fleet-mail queues a message in the target session inbox', async ($, on) => {
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW)
const answer = await $.command.run({
command: 'fleet-mail',
args: PEER + ' rebasing on main is safe now',
})
expect(answer.text).toContain('Queued for ' + PEER)
const inbox = store.get('inbox:' + PEER)
expect(inbox.length).toBe(1)
expect(inbox[0].text).toBe('rebasing on main is safe now')
expect(inbox[0].from).toBe(SELF)
})
test('a second message appends rather than replacing the first', async ($, on) => {
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW)
await $.command.run({ command: 'fleet-mail', args: PEER + ' first' })
await $.command.run({ command: 'fleet-mail', args: PEER + ' second' })
const inbox = store.get('inbox:' + PEER)
expect(inbox.length).toBe(2)
expect(inbox[0].text).toBe('first')
expect(inbox[1].text).toBe('second')
})
test('/fleet-mail refuses a target that never announced itself', async ($, on) => {
const store = stubEngine(on)
const answer = await $.command.run({ command: 'fleet-mail', args: 'ghost-session hello' })
expect(answer.text).toContain('No live fleet session ghost-session')
// Nothing may be queued for a session we could not confirm.
expect(store.get('inbox:ghost-session')).toBe(undefined)
})
test('/fleet-mail rejects input with no message text', async ($, on) => {
const store = stubEngine(on)
announce(store, PEER, '/Users/x/work-repo', NOW)
for (const args of ['', PEER, PEER + ' ']) {
const answer = await $.command.run({ command: 'fleet-mail', args })
expect(answer.text).toContain('Usage: /fleet-mail')
}
expect(store.get('inbox:' + PEER)).toBe(undefined)
})
test('an inbound peer message is passed on, not consumed', async ($, on) => {
// A receive hook that withheld a message would silently break Claude Code's
// own channel, so assert this mod stays transparent.
stubEngine(on)
const result = await $.session.receive({
text: 'from the other session',
origin: { kind: 'peer' },
})
expect(result?.consumed).toBe(undefined)
})
/** Answer fleetd's three MCP requests the way the live daemon does. */
function stubFleetd(on: any, toolText: string, seen: any[]) {
on('http.fetch', (_$: any, e: any) => {
const body = JSON.parse(e.init.body)
seen.push({ method: body.method, session: e.init.headers['Mcp-Session-Id'] })
if (body.method === 'initialize') {
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-1' }, text: '{}' } }
}
if (body.method === 'tools/call') {
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text: toolText }] } }
return { value: { ok: true, status: 200, headers: {}, text: 'event: message\ndata: ' + JSON.stringify(result) + '\n' } }
}
return { value: { ok: true, status: 202, headers: {}, text: '' } }
})
}
test('/fleet-whoami reads the tool result out of the event stream', async ($, on) => {
stubEngine(on)
const seen: any[] = []
stubFleetd(on, '{"role":"observer","sessionId":"term_x"}', seen)
const answer = await $.command.run({ command: 'fleet-whoami', args: '' })
expect(answer.text).toBe('fleetd says: {"role":"observer","sessionId":"term_x"}')
// The session id from initialize must ride on every later request.
expect(seen.map((r) => r.method)).toEqual(['initialize', 'notifications/initialized', 'tools/call'])
expect(seen[2].session).toBe('sid-1')
})
test('/fleet-whoami reports a daemon that is down instead of throwing', async ($, on) => {
stubEngine(on)
on('http.fetch', () => ({ value: { ok: false, status: 503, headers: {}, text: '' } }))
const answer = await $.command.run({ command: 'fleet-whoami', args: '' })
expect(answer.text).toBe('fleetd unreachable: fleetd initialize failed with status 503')
})
/**
* The session.start hook's own calls, for a test that fires it. Kept apart from stubEngine
* because a timer test drives $.clock through mock.clock(on) instead of a fixed clock.now.
*/
function stubSessionStart(on: any, submitted: string[]): Map<string, any> {
const store = new Map<string, any>()
on('store.get', (_$: any, e: any) => ({ value: store.get(e.key) }))
on('store.set', (_$: any, e: any) => {
store.set(e.key, e.value)
return { value: undefined }
})
on('store.keys', () => ({ value: [...store.keys()] }))
on('session.id', () => ({ value: SELF }))
on('session.cwd', () => ({ value: '/Users/x/claude-bridge' }))
on('command.register', () => ({ value: undefined }))
on('ui.log', () => ({ value: undefined }))
on('env.get', () => ({ value: undefined }))
// The engine skips a prompt.submit hook that answers anything but { text } or { drop }, and
// the mod's callback then throws, so this must hand the text straight back.
on('prompt.submit', (_$: any, e: any) => {
submitted.push(e.text)
return { text: e.text }
})
on('session.start', (_$: any, e: any) => e)
return store
}
/**
* Answer fleetd's MCP requests, with the tool result read fresh on every call.
*
* `fetches` collects every method sent, with a tools/call entry naming its tool (such as
* `tools/call:fleet_inbox`), and `sessions` the Mcp-Session-Id of each tools/call. Each
* initialize hands out the next id, so a reused session and a reopened one differ.
* `sessionGone` makes a tools/call answer the way fleetd answers for a session id it no
* longer holds. `whoamiText` answers a `fleet_whoami` call apart from `toolText`, which
* answers every other tool.
*/
function stubFleetdDynamic(
on: any,
toolText: () => string,
fetches: string[],
sessions: string[] = [],
sessionGone: () => boolean = () => false,
whoamiText: () => string = () => '{"role":"primary"}',
) {
let opened = 0
on('http.fetch', (_$: any, e: any) => {
const body = JSON.parse(e.init.body)
if (body.method === 'initialize') {
fetches.push(body.method)
opened += 1
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-' + opened }, text: '{}' } }
}
if (body.method === 'tools/call') {
fetches.push(body.method + ':' + body.params.name)
sessions.push(e.init.headers['Mcp-Session-Id'])
if (sessionGone()) {
const gone = '{"jsonRpcError":{"code":-32603,"message":"Session not found"}}'
return { value: { ok: false, status: 404, headers: {}, text: gone } }
}
const text = body.params.name === 'fleet_whoami' ? whoamiText() : toolText()
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text }] } }
return { value: { ok: true, status: 200, headers: {}, text: 'data: ' + JSON.stringify(result) + '\n' } }
}
fetches.push(body.method)
return { value: { ok: true, status: 202, headers: {}, text: '' } }
})
}
/** How many of `fetches` were an initialize. */
function initializes(fetches: string[]): number {
return fetches.filter((m) => m === 'initialize').length
}
test('the fleetd inbox poll submits each collected message and names the sender', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
let inbox = { sessionId: 'term_self', count: 0, messages: [] as string[] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), [])
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
// An empty inbox is the ordinary answer and must submit nothing.
await clock.advance(3_000)
expect(submitted.length).toBe(0)
// The control for the line above: the same timer, the same stubs, with mail waiting.
inbox = { sessionId: 'term_self', count: 2, messages: ['rebase is safe now', 'build is green'] }
await clock.advance(3_000)
expect(submitted.length).toBe(2)
expect(submitted[0]).toContain('rebase is safe now')
expect(submitted[0]).toContain('fleetd')
expect(submitted[1]).toContain('build is green')
})
test('a collected message is not submitted a second time', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
// fleetd removes a message when it hands it over, so the next poll answers empty.
let inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), [])
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
expect(submitted.length).toBe(1) // control: the first poll really did deliver it
inbox = { sessionId: 'term_self', count: 0, messages: [] }
await clock.advance(3_000)
expect(submitted.length).toBe(1)
})
test('a fleetd that is down leaves the poll timer running', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
on('http.fetch', (_$: any, e: any) => {
fetches.push(JSON.parse(e.init.body).method)
return { value: { ok: false, status: 503, headers: {}, text: '' } }
})
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
const afterFirst = fetches.length
expect(afterFirst).toBeGreaterThan(0)
expect(submitted.length).toBe(0)
// A throw out of the timer would stop it, so a second tick that still reaches fleetd is the
// positive control for the assertion above: the timer survived the failure.
await clock.advance(3_000)
expect(fetches.length).toBeGreaterThan(afterFirst)
expect(submitted.length).toBe(0)
})
test('a second inbox poll reuses the first MCP session', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const sessions: string[] = []
const inbox = { sessionId: 'term_self', count: 0, messages: [] as string[] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, sessions)
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
// The first poll also reads the role once, so it makes two tools/call (whoami, then
// inbox); the second poll already knows the role and makes only one (inbox).
expect(sessions.length).toBe(3) // control: all three calls really reached fleetd
expect(initializes(fetches)).toBe(1)
expect(sessions.every((s) => s === sessions[0])).toBe(true)
})
test('a session fleetd no longer holds is opened again and the call retried', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const sessions: string[] = []
let gone = false
const inbox = { sessionId: 'term_self', count: 1, messages: ['the daemon restarted'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, sessions, () => {
// Only the first call after the flag is set is refused; the retry succeeds.
const refuse = gone
gone = false
return refuse
})
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
// The control: an ordinary poll opens one session and needs no second one.
await clock.advance(3_000)
expect(initializes(fetches)).toBe(1)
expect(submitted.length).toBe(1)
gone = true
await clock.advance(3_000)
expect(initializes(fetches)).toBe(2)
expect(sessions[sessions.length - 1]).toBe('sid-2')
expect(submitted.length).toBe(2)
})
/** How many of `fetches` were a tools/call for the named tool. */
function toolCalls(fetches: string[], tool: string): number {
return fetches.filter((m) => m === 'tools/call:' + tool).length
}
test('a worker role stops the poll from calling fleet_inbox', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"worker"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_whoami')).toBe(1)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(0)
expect(submitted.length).toBe(0)
})
test('an architect role stops the poll from calling fleet_inbox', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"architect"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_whoami')).toBe(1)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(0)
expect(submitted.length).toBe(0)
})
test('a primary role keeps the poll calling fleet_inbox', async ($, on) => {
// The positive control for the two tests above: the same stubs, a role neither
// gates, so a bug that silenced every role would still pass them.
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"primary"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(1)
expect(submitted.length).toBe(1)
})
test('an observer role keeps the poll calling fleet_inbox', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"observer"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(1)
expect(submitted.length).toBe(1)
})
test('a whoami that fails on the first tick is retried and the poll resumes', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
let whoamiFails = true
on('http.fetch', (_$: any, e: any) => {
const body = JSON.parse(e.init.body)
if (body.method === 'tools/call' && body.params.name === 'fleet_whoami' && whoamiFails) {
fetches.push(body.method + ':' + body.params.name)
return { value: { ok: false, status: 503, headers: {}, text: '' } }
}
if (body.method === 'initialize') {
fetches.push(body.method)
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-1' }, text: '{}' } }
}
if (body.method === 'tools/call') {
fetches.push(body.method + ':' + body.params.name)
const text = body.params.name === 'fleet_whoami' ? '{"role":"observer"}' : JSON.stringify(inbox)
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text }] } }
return { value: { ok: true, status: 200, headers: {}, text: 'data: ' + JSON.stringify(result) + '\n' } }
}
fetches.push(body.method)
return { value: { ok: true, status: 202, headers: {}, text: '' } }
})
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(0) // control: a failed whoami calls no fleet_inbox
expect(submitted.length).toBe(0)
whoamiFails = false
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(1)
expect(submitted.length).toBe(1)
})
test('whoami is called once, not on every tick, once the role is known', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 0, messages: [] as string[] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"primary"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_whoami')).toBe(1)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(3)
})