Compare commits
16 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6022aef612 | |||
| 66ab9cab24 | |||
| f611488c9b | |||
| 44d16e7639 | |||
| 87f7c03535 | |||
| 562354a58d | |||
| eaa1f0d170 | |||
| 72ea7def0a | |||
| 98f3cd5aa7 | |||
| aec5ff2c2f | |||
| d9fedfbc5a | |||
| ac535503ff | |||
| 654b3b5e14 | |||
| 1fae8b81a0 | |||
| 4ef815682b | |||
| 6596458ce6 |
@@ -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"
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
},
|
||||
|
||||
@@ -1,8 +0,0 @@
|
||||
{
|
||||
"mcpServers": {
|
||||
"fleet": {
|
||||
"type": "http",
|
||||
"url": "${FLEETD_MCP_URL}"
|
||||
}
|
||||
}
|
||||
}
|
||||
+35
-19
@@ -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
|
||||
|
||||
|
||||
@@ -0,0 +1,4 @@
|
||||
{
|
||||
"description": "fleet mod: cross-account session messaging and PTY-free delivery",
|
||||
"modules": ["./register.js"]
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
}
|
||||
@@ -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:**
|
||||
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
Reference in New Issue
Block a user