Compare commits
28 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7dad054045 | |||
| 4e414c475f | |||
| 9084667493 | |||
| 2289e94223 | |||
| 92adfcfae5 | |||
| 5f7f388e69 | |||
| 39accf73e6 | |||
| 5c2f296bc3 | |||
| 0f2ec7a6b5 | |||
| 02eff4c532 | |||
| 92b6d9406b | |||
| c4498607e1 | |||
| e42eab5b4c | |||
| b1d2cb48ac | |||
| 001367d82c | |||
| 9fdcaaa8fd | |||
| 459a523e2c | |||
| 291dc02c77 | |||
| be835aa259 | |||
| 80506c79d0 | |||
| 6677ec8c63 | |||
| ca0c965932 | |||
| e8ab933cbc | |||
| 2adb950a12 | |||
| 7f0c4a8464 | |||
| 682991a846 | |||
| abe617c48c | |||
| 8a1d73b39e |
@@ -133,8 +133,17 @@ A handover file is a record of state and decisions. It is not a diary.
|
||||
**Run the three steps in this order. The order is not a style choice — the wrong order is
|
||||
refused.**
|
||||
|
||||
1. **`fleet_handover{action: "open", reason: "<why now>"}`.** It returns a `token` and the
|
||||
`handoverPath` you must write to. Nothing has happened to your pane yet.
|
||||
1. **`fleet_handover{action: "open", reason: "<why now>"}`.** It returns a `token`, the
|
||||
`handoverPath` you must write to, and `outstandingTickets` plus `openAsks`. Nothing has happened
|
||||
to your pane yet.
|
||||
|
||||
**Copy `outstandingTickets` and `openAsks` into the handover file.** Your successor keeps the
|
||||
authority to poll those tickets and answer those asks, because both are gated on the lead's
|
||||
name, which does not change when your pane does. What it does not keep is the ids — they exist
|
||||
only in your context and in this response. A ticket already in a terminal phase is the urgent
|
||||
one: its reply lives only in memory and is deleted once the ticket TTL passes, so an uncollected
|
||||
report is lost for good. Poll those before you confirm, or name them in the file so your
|
||||
successor polls them first.
|
||||
|
||||
**Write to exactly that path, and do not resolve it yourself.** It is always absolute, even when
|
||||
the operator configured a relative `handoverPath`: fleetd resolves a relative one against your
|
||||
|
||||
@@ -33,8 +33,11 @@ gate uses. A worker also carries its `sessionId`, `profile`, `worktree` and `bra
|
||||
carries the slot name it was bound to; a collaborator carries its registry name and its own
|
||||
`sessionId`, and **no `leader` key** — a collaborator is a named peer, not a primary. An **observer**
|
||||
carries only its own `sessionId`: a pane the daemon could not place as any of the above, authorized
|
||||
to `READ`/`METRICS` and to `REPLY`/`ASK` on its own pane and nothing more — never `SEND`, never a
|
||||
ticket. Don't infer what you can ask.
|
||||
to `READ`/`METRICS`, to `REPLY`/`ASK` on its own pane, and to `SEND` only to a target that resolves
|
||||
as an observer too — never to a lead, a collaborator, or a spawned member, and never a ticket. It
|
||||
finds such a target in `fleet_list`'s `panes` array, which for an observer is filtered to exactly
|
||||
what it may send to and reduced to `sessionId`, `label`, `status`, `role` and `deliverable`.
|
||||
Don't infer what you can ask.
|
||||
|
||||
Only if that call is unavailable, fall back to these — each is one-way, so keep reading until one
|
||||
fires: the reply charter in your system prompt (*"You are a spawned member in the
|
||||
@@ -62,9 +65,10 @@ and the sender silently receives nothing. Fail toward the recoverable error.
|
||||
2. **The bridge is the only channel.** Text you print in your terminal reaches nobody — the other
|
||||
side cannot see your screen. An answer that isn't in a `fleet_*` call is silently discarded.
|
||||
3. **Identity comes from the connection, never an argument.** Workers never pass a target; you
|
||||
cannot act as another session. Spawn/stop/drain are lead-only; **send is lead, architect, or
|
||||
collaborator** — and a collaborator may send only to a lead or another collaborator, never to a
|
||||
spawned member's terminal; reply/ask are only-as-itself — any peer may answer for its own pane,
|
||||
cannot act as another session. Spawn/stop/drain are lead-only; **send is lead, architect,
|
||||
collaborator, or observer** — and a collaborator may send only to a lead or another collaborator,
|
||||
never to a spawned member's terminal, while an observer may send only to another observer pane;
|
||||
reply/ask are only-as-itself — any peer may answer for its own pane,
|
||||
and for no other. A call outside your role is refused, not queued.
|
||||
4. **Delivery is status-gated: one message per turn.** Don't busy-poll a peer's terminal and don't
|
||||
re-send because a call looks slow — the bridge delivers when the peer is `idle`, `blocked` or
|
||||
@@ -178,6 +182,7 @@ you decide.
|
||||
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
|
||||
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
|
||||
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` reports a `collaborators` array, and each row carries that peer's `name` and the `sessionId` you send to. It is visible to you, to an architect and to another collaborator, never to a worker. Coordination only, **never** a task |
|
||||
| Message an **unconfigured pane** — a tab a person opened by hand | `fleet_send{sessionId: <their terminal>, content}` — it needs **no** `fleet.collaborators` entry and no restart, because a pane becomes deliverable the moment its agent connects the bridge MCP. `fleet_list`'s `panes` array reports every such pane with its label and the terminal id to send to — the full row for you, an architect or a collaborator; filtered and reduced for an observer. **`ListAgents` still never lists these**, and joining `herdr tab list` to `GET /agents` on `tab_id` stays the read-only fallback if the array is missing. Such a pane resolves as an `observer`: it can answer you with `fleet_reply`, and it can `fleet_send` to another observer pane, but never to you. Coordination only, **never** a task |
|
||||
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
|
||||
| Read your own held lead-to-lead mail (no ack) | `fleet_poll{coordId: <your own coord-id, from fleet_list's coordinator.selfId>}` — primary-only; never acks, so `fleet_list`'s `held[]` still shows it after. `fleet_list`'s `held[]` gives only a truncated preview — this is the only way to read the full body |
|
||||
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
|
||||
@@ -502,3 +507,33 @@ to replace them.
|
||||
|
||||
Prefer the unnamed lambda parameter `_` for required-but-unused params; a non-public
|
||||
`static void main(String[])` is valid (JEP 512) and boots via `java -jar`.
|
||||
|
||||
## Code quality — five rules, and what each already cost (enforced)
|
||||
|
||||
Measured at `7f0c4a8`: 124 main files, 36,278 lines, of which **17,850 are code** — 44% is comment,
|
||||
and only **two** files exceed 1000 *code* lines. Encapsulation and inheritance are already sound (0
|
||||
public mutable fields; 12 `extends`, 8 of them exceptions; 120 records). So there is **no Clean Code
|
||||
section, no SOLID list and no pattern catalogue** here: two architect reviews rejected those
|
||||
independently as text that would change no behaviour. These five rules are the whole standard.
|
||||
|
||||
1. **A comment states the current contract or a current maintainer constraint — nothing else.** No
|
||||
tickets, history, dates, measurements or review rationale; those go in the commit message or the
|
||||
MR description. Source code only — this rule never applies to Markdown.
|
||||
2. **A javadoc block stops at 30 lines.** Longer means it is a design argument, so it moves to
|
||||
`docs/<subject>.md` and is linked in one line. The longest here is 235 lines
|
||||
(`config/ConfigRef.java`) and the knowledge in it is load-bearing: **move it, never delete it.**
|
||||
This project has **no ADR** — subject pages under `docs/` are the destination.
|
||||
3. **A comment in main source never names a test class.** There are 44 such names in 76 places and
|
||||
**2 are already dead**, because a name inside `{@code}` is invisible to the compiler and rots in
|
||||
silence. Say what the code guarantees; the test is found by looking.
|
||||
4. **Never relieve a testing problem by reshaping production code.** `Fleetd` carries 54 static
|
||||
factories, `MessageService` carries 7 `volatile` race hooks, and 8 tests assert on main source as
|
||||
*text*. Make the part injectable instead. `FleetdAssembly.assembleAndStart` is 452 lines and may
|
||||
not grow; no new source-text test may be added.
|
||||
5. **No new package cycle, and no widening of a recorded one.** Five pairs are frozen as an exact
|
||||
edge baseline in `PackageCyclesTest` — four of them involve `msg`.
|
||||
|
||||
Rules 1, 2, 3 and 5 have build checks, and Gitea CI runs them on every PR, so they bind members too.
|
||||
**Rule 4's judgement half has no mechanism**: a cap stops a count growing, but no test tells a good
|
||||
decomposition from a bad one. That half is a review obligation, and saying so is deliberate — a rule
|
||||
dressed as a gate it does not have is worse than an honest review item.
|
||||
|
||||
@@ -182,6 +182,12 @@ Two consequences a lead feels directly:
|
||||
is fine; the message simply waits, and then restarts the member when it next goes idle.
|
||||
- **A spawned member is not deliverable until it has mounted the MCP.** Until then a send waits on
|
||||
that gate for about 60 seconds and then fails without ever reaching the pane.
|
||||
- **The same gate is what makes an unconfigured pane deliverable.** `contextExtractor` runs on
|
||||
every MCP request, `initialize` included, and `markTrackedCallerPresent` enrols a spawned member
|
||||
*or* an observer into `MemberPresence`; `deliverableTo` then tests presence before the lead and
|
||||
collaborator maps. So mounting the server is the enrolment, and a tab a person opened by hand can
|
||||
be sent to with no config and no restart. It answers with `fleet_reply` — it cannot `fleet_send`,
|
||||
because `Authz` keeps `SEND` to a primary, an architect or a collaborator.
|
||||
|
||||
`UNKNOWN` is deliberately neither injectable nor a pickup. A pane whose status cannot be read is
|
||||
not a pane that is safe to write to — see fleetd #176 for what happens when a gate treats an
|
||||
|
||||
@@ -234,7 +234,7 @@ public final class Fleetd {
|
||||
* <p>Both sets are read through their supplier on each call rather than snapshotted, so a lead or
|
||||
* collaborator discovered by {@code leadScan} after startup becomes deliverable without a restart.
|
||||
*/
|
||||
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads,
|
||||
public static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads,
|
||||
Supplier<Map<String, String>> collaborators) {
|
||||
return target -> presence.isPresent(target) || leads.get().containsKey(target)
|
||||
|| collaborators.get().containsKey(target);
|
||||
@@ -336,22 +336,6 @@ public final class Fleetd {
|
||||
}, reasonByCredential::get);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #415 (review follow-up): package-private factory for the CB-578 stage A {@code
|
||||
* exhaustedPattern} startup coverage line, paired explicitly with {@link
|
||||
* CompletionResolver.UnsetMeaning#OFF} — {@code exhaustedPattern} has no fallback, so a
|
||||
* profile with none configured really does have the classification off.
|
||||
*
|
||||
* <p>Extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
|
||||
* #worktreeBranchLookup} were: {@code coverage()}'s own tests ({@code CompletionResolverTest})
|
||||
* prove it words {@code OFF} and {@link CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT}
|
||||
* correctly when a test supplies the meaning itself — they cannot prove {@code main} pairs the
|
||||
* right meaning with the right key, which is the actual fleetd #415 defect. <b>Measured:</b>
|
||||
* swapping the {@code UnsetMeaning} arguments between this method and {@link
|
||||
* #errorPatternCoverageLine} — recreating #415's defect with the two keys exchanged — compiled
|
||||
* with 0 errors and left all 1506 existing tests green before {@code
|
||||
* FleetdPatternCoverageLineTest} was added to catch exactly that swap.
|
||||
*/
|
||||
/**
|
||||
* fleetd #446 follow-up: the criterion-2 WARNING text — "name the fix, not just the fact" —
|
||||
* for a profile whose {@code model:} is configured. Extracted out of the {@code
|
||||
@@ -386,19 +370,22 @@ public final class Fleetd {
|
||||
+ "s quarantine above is the only thing keeping new spawns off it for now";
|
||||
}
|
||||
|
||||
/**
|
||||
* Package-private factory for the {@code exhaustedPattern} startup coverage line, paired
|
||||
* explicitly with {@link CompletionResolver.UnsetMeaning#OFF} — {@code exhaustedPattern} has
|
||||
* no fallback, so a profile with none configured really does have the classification off.
|
||||
*/
|
||||
static String exhaustedPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
return CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
|
||||
allProfiles, configuredProfiles);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #415 (review follow-up): the {@code errorPattern} counterpart of {@link
|
||||
* #exhaustedPatternCoverageLine}, paired explicitly with {@link
|
||||
* CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT} — an unset {@code errorPattern} still runs
|
||||
* backend-error classification against {@code CompletionResolver}'s built-in {@code
|
||||
* BACKEND_ERROR} pattern, so the empty case is not "off". See {@link
|
||||
* #exhaustedPatternCoverageLine}'s javadoc for the measured swap mutation this pairing guards
|
||||
* against.
|
||||
* The {@code errorPattern} counterpart of {@link #exhaustedPatternCoverageLine}, paired
|
||||
* explicitly with {@link CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT} — an unset
|
||||
* {@code errorPattern} still runs backend-error classification against
|
||||
* {@code CompletionResolver}'s built-in {@code BACKEND_ERROR} pattern, so the empty case is
|
||||
* not "off".
|
||||
*/
|
||||
static String errorPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
return CompletionResolver.coverage("errorPattern", CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT,
|
||||
@@ -734,8 +721,8 @@ public final class Fleetd {
|
||||
* {@code CompletionResolver} constructor call — provably untested wiring, the whole reason
|
||||
* fleetd #248 exists: dropping that one argument (passing {@code _ -> null} instead) compiled
|
||||
* clean and left every test green. Extracted here, {@code main} now calls this factory instead
|
||||
* of building the lambda inline, and a source assertion on that call site
|
||||
* ({@code FleetdCompletionResolverWiringTest}) proves the argument is still actually passed.
|
||||
* of building the lambda inline, so the argument reaching the {@code CompletionResolver}
|
||||
* constructor is a named, directly testable call rather than an inline lambda.
|
||||
*
|
||||
* <p>Takes the roster as a plain {@link Supplier} — not a {@link SessionManager} — so this is
|
||||
* directly testable with a hand-built session list; no real {@code SessionManager} (launcher,
|
||||
|
||||
@@ -70,33 +70,58 @@ public final class Authz {
|
||||
*/
|
||||
public static final Predicate<String> NO_KNOWN_LEAD_OR_COLLABORATOR = target -> false;
|
||||
|
||||
/**
|
||||
* The fail-closed classifier for an observer's {@code SEND}: answers no for every target, so
|
||||
* the grant is refused unless a caller supplies a real one. {@code
|
||||
* CallerResolver#sendableObserverTarget()} is the real one, read from the same maps {@code
|
||||
* CallerResolver#resolve} consults, so a target that classifier calls known is one {@code
|
||||
* resolve} would actually resolve as {@link Role#OBSERVER}.
|
||||
*/
|
||||
public static final Predicate<String> NO_KNOWN_OBSERVER_TARGET = target -> false;
|
||||
|
||||
/**
|
||||
* Convenience form for a caller with no classifier to supply. Fails closed: a collaborator's
|
||||
* {@code SEND} is refused, as if no terminal were a configured lead or collaborator — the
|
||||
* same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR} gives explicitly. Every other action's
|
||||
* result is identical to the four-argument form's, since none of them consult the classifier.
|
||||
* or an observer's {@code SEND} is refused, as if no terminal were a configured lead,
|
||||
* collaborator, or observer target — the same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR}
|
||||
* and {@link #NO_KNOWN_OBSERVER_TARGET} give explicitly. Every other action's result is
|
||||
* identical to the five-argument form's, since none of them consult either classifier.
|
||||
*
|
||||
* <p>Its default classifier denies every collaborator, so a caller enforcing authorization
|
||||
* must use the four-argument form instead.
|
||||
* <p>Its default classifiers deny every collaborator and every observer, so a caller
|
||||
* enforcing authorization must use the five-argument form instead.
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession) {
|
||||
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR);
|
||||
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR, NO_KNOWN_OBSERVER_TARGET);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #permits(Principal, Action, String)}, with a real classifier for a collaborator's
|
||||
* {@code SEND}. An observer's {@code SEND} still fails closed ({@link #NO_KNOWN_OBSERVER_TARGET}) —
|
||||
* a caller enforcing both grants must use the five-argument form.
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession,
|
||||
Predicate<String> knownLeadOrCollaborator) {
|
||||
return permits(caller, action, targetSession, knownLeadOrCollaborator, NO_KNOWN_OBSERVER_TARGET);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code caller} may perform {@code action} against {@code targetSession}.
|
||||
*
|
||||
* @param targetSession the session id in the request path; only consulted for the
|
||||
* worker-scoped actions ({@code REPLY}, {@code ASK}) and for a
|
||||
* collaborator's {@code SEND}, ignored otherwise, may be
|
||||
* {@code null}
|
||||
* worker-scoped actions ({@code REPLY}, {@code ASK}), for a
|
||||
* collaborator's {@code SEND}, and for an observer's
|
||||
* {@code SEND}, ignored otherwise, may be {@code null}
|
||||
* @param knownLeadOrCollaborator whether a terminal is a configured lead or collaborator —
|
||||
* consulted only for a collaborator's {@code SEND}, to confine
|
||||
* it to another named peer and never a spawned member's
|
||||
* terminal
|
||||
* @param knownObserverTarget whether a terminal is one this daemon would itself resolve as
|
||||
* {@link Role#OBSERVER} — consulted only for an observer's
|
||||
* {@code SEND}, to confine it to another observer pane and never
|
||||
* a lead, a collaborator, or a spawned member
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession,
|
||||
Predicate<String> knownLeadOrCollaborator) {
|
||||
Predicate<String> knownLeadOrCollaborator,
|
||||
Predicate<String> knownObserverTarget) {
|
||||
if (caller == null || caller.isAnonymous()) {
|
||||
return false; // authenticated as nothing ⇒ authorized for nothing
|
||||
}
|
||||
@@ -107,13 +132,15 @@ public final class Authz {
|
||||
// would be a worker escalating into the orchestrator role.
|
||||
case SPAWN, STOP, DRAIN, HANDOVER -> caller.isPrimary();
|
||||
|
||||
// Delivering a turn to a local session is open to the primary, the architect, and a
|
||||
// collaborator whose target is itself a configured lead or collaborator: the architect
|
||||
// delegates to workers (that is the role's point); a collaborator may reach only
|
||||
// another named peer, never a spawned member's terminal. A worker is excluded —
|
||||
// sending would be it escalating.
|
||||
// Delivering a turn to a local session is open to the primary and the architect
|
||||
// unconditionally. A collaborator may reach only a target that is itself a configured
|
||||
// lead or collaborator, never a spawned member's terminal. An observer may reach only
|
||||
// a target that would itself resolve as an observer, never a lead, a collaborator, or
|
||||
// a spawned member. A worker is excluded from every case — sending would be it
|
||||
// escalating into the orchestrator role.
|
||||
case SEND -> caller.isPrimary() || caller.isArchitect()
|
||||
|| (caller.isCollaborator() && knownLeadOrCollaborator.test(targetSession));
|
||||
|| (caller.isCollaborator() && knownLeadOrCollaborator.test(targetSession))
|
||||
|| (caller.isObserver() && knownObserverTarget.test(targetSession));
|
||||
|
||||
// Resolving a worker's blocked question is part of delegating to it, open to the same
|
||||
// two roles that may stand up that delegation in the first place. Not a collaborator:
|
||||
|
||||
@@ -270,6 +270,31 @@ public final class CallerResolver {
|
||||
|| collaboratorTerminals.get().containsKey(target);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code target} names a terminal this resolver would itself resolve as {@link
|
||||
* Role#OBSERVER} — the classifier an observer's {@code SEND} is checked against, read from the
|
||||
* same maps and functions {@link #resolve} consults so a target this accepts is exactly one
|
||||
* {@code resolve} would hand back {@link Role#OBSERVER} for, and the reverse.
|
||||
*/
|
||||
public Predicate<String> sendableObserverTarget() {
|
||||
return target -> target != null
|
||||
&& spawnedMemberRole.apply(target) == null
|
||||
&& !leadTerminals.get().containsKey(target)
|
||||
&& !boundToArchitectSlot(target)
|
||||
&& !collaboratorTerminals.get().containsKey(target);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code terminal} is bound to a configured slot the live roster still confirms as an
|
||||
* architect — the one classifier {@link #sendableObserverTarget()} and {@code FleetMcp}'s
|
||||
* {@code panes} row both read, so a pane's reported role and its {@code SEND} reachability can
|
||||
* never drift apart.
|
||||
*/
|
||||
public boolean boundToArchitectSlot(String terminal) {
|
||||
String slot = architectTerminals.get().get(terminal);
|
||||
return slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT;
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve the caller of a request.
|
||||
*
|
||||
|
||||
@@ -12,9 +12,10 @@ package dev.ltms.fleet.auth;
|
||||
public enum Role {
|
||||
|
||||
/**
|
||||
* The orchestrating session. Established either by being a loopback caller that is not a
|
||||
* worker pane (under {@code loopback-trust}) or by presenting a valid bearer token (under
|
||||
* {@code token} mode).
|
||||
* The orchestrating session. Established either by being a loopback caller that resolves to
|
||||
* no herdr pane at all (under {@code loopback-trust}) or by presenting a valid bearer token
|
||||
* (under {@code token} mode). A loopback caller that does own a pane, but matches none of the
|
||||
* roles below, resolves to {@link #OBSERVER} instead.
|
||||
*/
|
||||
PRIMARY,
|
||||
|
||||
@@ -49,10 +50,11 @@ public enum Role {
|
||||
* A loopback pane that resolved to none of the roles above: not a live spawned member, not a
|
||||
* configured lead, not a bound architect slot, not a configured collaborator tab. Unforgeable
|
||||
* like a worker's — derived from the connection's pane, never from a request argument, and
|
||||
* honoured regardless of auth mode. May {@code READ} and {@code METRICS}, and {@code REPLY}/
|
||||
* {@code ASK} only as its own pane; may not {@code SPAWN}/{@code STOP}/{@code DRAIN}/
|
||||
* {@code HANDOVER}, {@code SEND}, poll a ticket ({@code TASK_READ}), or reach the coordination
|
||||
* broker ({@code COORD_SEND}/{@code COORD_READ}).
|
||||
* honoured regardless of auth mode. May {@code READ} and {@code METRICS}, {@code REPLY}/
|
||||
* {@code ASK} only as its own pane, and {@code SEND} only to a target that would itself
|
||||
* resolve as {@code OBSERVER}; may not {@code SPAWN}/{@code STOP}/{@code DRAIN}/
|
||||
* {@code HANDOVER}, poll a ticket ({@code TASK_READ}), or reach the coordination broker
|
||||
* ({@code COORD_SEND}/{@code COORD_READ}).
|
||||
*/
|
||||
OBSERVER,
|
||||
|
||||
|
||||
@@ -2095,13 +2095,6 @@ public record FleetConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The top-level keys in {@code yaml} that this build does not understand, sorted. Package-private
|
||||
* so the guardrail is asserted directly rather than through a log appender.
|
||||
*
|
||||
* @return empty when everything is known, or when {@code yaml} is not a mapping at all (a
|
||||
* malformed file is {@code readValue}'s error to report, not this method's)
|
||||
*/
|
||||
/**
|
||||
* Top-level keys renamed by the member taxonomy, mapped old → new.
|
||||
*
|
||||
@@ -2636,6 +2629,13 @@ public record FleetConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The top-level keys in {@code yaml} that this build does not understand, sorted. Package-private
|
||||
* so the guardrail is asserted directly rather than through a log appender.
|
||||
*
|
||||
* @return empty when everything is known, or when {@code yaml} is not a mapping at all (a
|
||||
* malformed file is {@code readValue}'s error to report, not this method's)
|
||||
*/
|
||||
static List<String> unknownTopLevelKeys(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
|
||||
@@ -4,6 +4,7 @@ import com.fasterxml.jackson.databind.JsonNode;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -145,6 +146,32 @@ public final class PaneLocator {
|
||||
return ancestry;
|
||||
}
|
||||
|
||||
/**
|
||||
* Every tab herdr tracks across every searched daemon, keyed by tab id, to its display label —
|
||||
* the pane-discovery surface behind {@code GET /agents} and {@code fleet_list}'s {@code panes}
|
||||
* row. Collapses to one scan in the single-daemon deployment, the same as
|
||||
* {@link #terminalForPid}. A tab herdr reports with no label maps to a {@code null} value here;
|
||||
* a tab with no {@code tab_id} is skipped.
|
||||
*/
|
||||
public Map<String, String> tabLabelsByTabId() {
|
||||
Map<String, String> out = new LinkedHashMap<>();
|
||||
for (HerdrClient herdr : herdrs) {
|
||||
for (JsonNode w : herdr.call("workspace.list").path("workspaces")) {
|
||||
String workspaceId = w.path("workspace_id").asText(null);
|
||||
if (workspaceId == null) {
|
||||
continue;
|
||||
}
|
||||
for (JsonNode t : herdr.call("tab.list", Map.of("workspace_id", workspaceId)).path("tabs")) {
|
||||
Tab tab = Tab.from(t);
|
||||
if (tab.tabId() != null) {
|
||||
out.put(tab.tabId(), tab.label());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Whether a pane owns one of the scanned pid's ancestors, or the check of it failed outright. */
|
||||
private enum Ownership { OWNS, DOES_NOT_OWN, UNKNOWN }
|
||||
|
||||
|
||||
@@ -0,0 +1,169 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
/**
|
||||
* Whether an agent pane's input box is clear for a delivery.
|
||||
*
|
||||
* <p>{@link AgentControl#send} pastes its text and submits it in the same call, so a delivery into a
|
||||
* pane whose input box already holds characters submits those characters too. {@link
|
||||
* AgentStatus#injectable()} cannot see that: it describes the agent, and an agent waiting at its
|
||||
* prompt reports the same status whether its box is empty or holds a half-typed line. This reads the
|
||||
* box itself.
|
||||
*
|
||||
* <p>Only a box that is positively empty clears the gate. A box with content, a pane this cannot
|
||||
* recognise, and a failed read all hold the delivery, because a held delivery is recoverable and a
|
||||
* submitted half-line is not. Every caller must therefore be a path that retries.
|
||||
*
|
||||
* <p>A pane that holds for {@link #HOLD_WARN_STREAK} consecutive checks gets one warning, so a box
|
||||
* that never clears is visible instead of silent. The warning repeats only after the box has cleared
|
||||
* again.
|
||||
*/
|
||||
public final class PromptBox {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(PromptBox.class);
|
||||
|
||||
/**
|
||||
* herdr {@code agent.read} source. {@code detection} is the region herdr itself uses for status
|
||||
* detection, so the input box is always drawn in it. It is not a short tail: it carries transcript
|
||||
* scrollback above the box, including earlier prompts the operator has already submitted, which is
|
||||
* why only the last box line on it is the live one.
|
||||
*/
|
||||
static final String PROBE_SOURCE = "detection";
|
||||
|
||||
/** Consecutive holds for one target before one warning is logged. */
|
||||
static final int HOLD_WARN_STREAK = 20;
|
||||
|
||||
/**
|
||||
* Input box markers, each matched only as a line's first characters: the caret the current TUI
|
||||
* draws, and the bordered box an older one drew. A marker further along a line is transcript text,
|
||||
* such as a caret inside something the operator quoted.
|
||||
*/
|
||||
private static final List<String> BOX_MARKERS = List.of("❯", "│ >");
|
||||
|
||||
/** Marker of a turn that is still generating; a box drawn under it is not a settled prompt. */
|
||||
private static final String ACTIVE_TURN_MARKER = "esc to interrupt";
|
||||
|
||||
/** Block glyphs a terminal capture can leave in an otherwise empty box for the cursor cell. */
|
||||
private static final String CURSOR_GLYPHS = "█▉▊▋▌▍▎▏";
|
||||
|
||||
/** What a box holds: nothing, unsubmitted characters, or a pane this cannot read as a box. */
|
||||
public enum State { EMPTY, DRAFT, UNREADABLE }
|
||||
|
||||
/** A box reading: its state, and how many characters it holds ({@code 0} unless {@code DRAFT}). */
|
||||
public record Reading(State state, int characters) {
|
||||
}
|
||||
|
||||
private final AgentControl agents;
|
||||
|
||||
/** Consecutive holds per target, so a box that never clears can be warned about once. */
|
||||
private final Map<String, Integer> holdStreaks = new ConcurrentHashMap<>();
|
||||
|
||||
public PromptBox(AgentControl agents) {
|
||||
this.agents = agents;
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code target}'s input box is empty, so a delivery would submit only its own text.
|
||||
* {@code false} means hold and come back; it never means the delivery failed.
|
||||
*/
|
||||
public boolean clearToSubmit(String target) {
|
||||
Reading reading = inspect(target);
|
||||
if (reading.state() == State.EMPTY) {
|
||||
holdStreaks.remove(target);
|
||||
return true;
|
||||
}
|
||||
int streak = holdStreaks.merge(target, 1, Integer::sum);
|
||||
if (streak == HOLD_WARN_STREAK) {
|
||||
log.warn("prompt box of {} has held a delivery {} times in a row ({}, {} character(s) in the box)"
|
||||
+ " — nothing is lost, delivery resumes once the box is empty",
|
||||
target, streak, reading.state(), reading.characters());
|
||||
} else {
|
||||
log.debug("prompt box of {} is {} ({} character(s)), holding delivery {}",
|
||||
target, reading.state(), reading.characters(), streak);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/** Read and classify {@code target}'s pane. A read failure reads as {@link State#UNREADABLE}. */
|
||||
private Reading inspect(String target) {
|
||||
String pane;
|
||||
try {
|
||||
pane = agents.read(target, PROBE_SOURCE);
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("prompt box read for {} failed, holding delivery: {}", target, e.getMessage());
|
||||
return new Reading(State.UNREADABLE, 0);
|
||||
}
|
||||
return classify(pane);
|
||||
}
|
||||
|
||||
/**
|
||||
* Classify a Claude Code TUI pane region. Pure, so it is unit-testable without herdr.
|
||||
*
|
||||
* <p>{@link State#EMPTY} needs two positive signals: the pane's last box line holds nothing after
|
||||
* its marker, and nothing below that line says a turn is still generating. Everything else is
|
||||
* {@link State#UNREADABLE} — a blank capture, or a region with no box line at all — so a pane this
|
||||
* does not understand holds the delivery rather than guessing it is safe.
|
||||
*
|
||||
* <p>The generating marker is looked for only from the box line down. Above it is scrollback, where
|
||||
* an earlier turn's marker survives; treating that as a live turn would make {@link State#EMPTY}
|
||||
* unreachable and hold every delivery forever.
|
||||
*
|
||||
* <p>Whitespace, a trailing box border and a cursor block count as nothing. Any other character
|
||||
* counts as the operator's unsubmitted text, including a placeholder hint a future TUI might draw
|
||||
* there; that direction holds a delivery it could have sent, which {@link #HOLD_WARN_STREAK} makes
|
||||
* visible.
|
||||
*/
|
||||
static Reading classify(String pane) {
|
||||
if (pane == null || pane.isBlank()) return new Reading(State.UNREADABLE, 0);
|
||||
int box = lastBoxLineStart(pane);
|
||||
if (box < 0) return new Reading(State.UNREADABLE, 0);
|
||||
String fromBox = pane.substring(box);
|
||||
if (fromBox.toLowerCase().contains(ACTIVE_TURN_MARKER)) return new Reading(State.UNREADABLE, 0);
|
||||
String content = boxContent(firstLine(fromBox));
|
||||
return content.isEmpty() ? new Reading(State.EMPTY, 0) : new Reading(State.DRAFT, content.length());
|
||||
}
|
||||
|
||||
/** Offset of the last line starting with a box marker, or {@code -1} if the region has none. */
|
||||
private static int lastBoxLineStart(String pane) {
|
||||
int found = -1;
|
||||
for (int start = 0; start <= pane.length(); ) {
|
||||
int end = pane.indexOf('\n', start);
|
||||
String line = pane.substring(start, end < 0 ? pane.length() : end);
|
||||
if (markerLength(line) > 0) found = start;
|
||||
if (end < 0) break;
|
||||
start = end + 1;
|
||||
}
|
||||
return found;
|
||||
}
|
||||
|
||||
/** Length of the box marker this line starts with, or {@code 0} if it starts with none. */
|
||||
private static int markerLength(String line) {
|
||||
for (String marker : BOX_MARKERS) {
|
||||
if (line.startsWith(marker)) return marker.length();
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
private static String firstLine(String text) {
|
||||
int newline = text.indexOf('\n');
|
||||
return newline < 0 ? text : text.substring(0, newline);
|
||||
}
|
||||
|
||||
/** The text the box holds: its own line after the marker, stripped of border, padding and cursor. */
|
||||
private static String boxContent(String boxLine) {
|
||||
String line = boxLine.substring(markerLength(boxLine)).stripTrailing();
|
||||
if (line.endsWith("│")) line = line.substring(0, line.length() - 1);
|
||||
StringBuilder content = new StringBuilder();
|
||||
for (char c : line.toCharArray()) {
|
||||
if (Character.isWhitespace(c) || CURSOR_GLYPHS.indexOf(c) >= 0) continue;
|
||||
content.append(c);
|
||||
}
|
||||
return content.toString();
|
||||
}
|
||||
}
|
||||
@@ -4,16 +4,17 @@ import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.Set;
|
||||
|
||||
/**
|
||||
* Tracks which workers are <em>available</em> — their Claude has booted and connected its MCP client
|
||||
* to the bridge (CB-113). This is the reliable readiness signal, unlike herdr's {@code agent_status},
|
||||
* which reports {@code idle} for a worker whose Claude is still booting. Delivering into that boot
|
||||
* window pastes into a not-yet-ready TUI (the text is lost) and wedges the worker's delivery state,
|
||||
* so the {@link Injector} holds the first delivery until the worker is present here.
|
||||
* Tracks which peers are <em>available</em> — their Claude has booted and connected its MCP
|
||||
* client to the bridge. For a spawned member this is the reliable readiness signal,
|
||||
* unlike herdr's {@code agent_status}, which reports {@code idle} while its Claude is still
|
||||
* booting. Delivering into that boot window pastes into a not-yet-ready TUI (the text is lost)
|
||||
* and wedges that member's delivery state, so the {@link Injector} holds a spawned member's
|
||||
* first delivery until it is present here.
|
||||
*
|
||||
* <p>Populated from the MCP transport: any MCP request whose connection resolves to a worker terminal
|
||||
* marks that worker present (its {@code initialize} is the first such contact). A worker that never
|
||||
* mounts the bridge MCP is never marked present — its sends stay queued until they time out, which is
|
||||
* correct (it could not have replied anyway).
|
||||
* <p>Populated from the MCP transport, for the peers whose deliverability rests on proving a live
|
||||
* MCP contact rather than on a configured registry entry. A peer that never mounts the bridge MCP
|
||||
* is never marked present — its sends stay queued until they time out, which is correct (it could
|
||||
* not have replied anyway).
|
||||
*/
|
||||
public class MemberPresence {
|
||||
|
||||
|
||||
@@ -149,9 +149,17 @@ public final class LeadRollover {
|
||||
* calling lead's workspace, before storing it here; see this class's
|
||||
* javadoc. This is the value the MCP layer hands back to the lead as
|
||||
* "write your file here", so callers may rely on it always being absolute.
|
||||
* @param rolloverKey the key {@link #confirm}'s single-flight claim is taken under. {@link
|
||||
* #open} resolves this once, here, from {@code leadTerminal} while the
|
||||
* calling lead is certainly still live: the lead's configured name, or
|
||||
* {@code leadTerminal} itself when no name resolves. Carried rather than
|
||||
* recomputed at release time, because by the time a roll's continuation
|
||||
* releases its claim the OLD terminal may no longer resolve to any name at
|
||||
* all — recomputing there would release a different key than the one the
|
||||
* claim was taken under.
|
||||
*/
|
||||
public record PendingRollover(String token, String leadTerminal, String handoverPath,
|
||||
long requestedAtMillis) {}
|
||||
long requestedAtMillis, String rolloverKey) {}
|
||||
|
||||
/** Which check refused a {@link #confirm} call, named so a caller can act on it. */
|
||||
public enum RefusalReason {
|
||||
@@ -176,9 +184,9 @@ public final class LeadRollover {
|
||||
*/
|
||||
HANDOVER_STALE,
|
||||
/**
|
||||
* This lead terminal already has a roll running: an earlier {@link #confirm} call claimed
|
||||
* it and that roll's continuation has not released it yet. {@code detail} names the lead
|
||||
* terminal and the token that holds the claim.
|
||||
* This lead already has a roll running: an earlier {@link #confirm} call claimed its
|
||||
* single-flight key (see {@link PendingRollover#rolloverKey}) and that roll's continuation
|
||||
* has not released it yet. {@code detail} names the key and the token that holds the claim.
|
||||
*/
|
||||
ROLL_ALREADY_RUNNING
|
||||
}
|
||||
@@ -340,9 +348,12 @@ public final class LeadRollover {
|
||||
private final Function<String, String> leadWorkspace;
|
||||
/**
|
||||
* Terminal id → that lead's configured name under {@code fleet.leaders}, or {@code null} when
|
||||
* the terminal names no currently-recognised lead. The deferred continuation calls this, on the
|
||||
* OLD terminal, before tearing it down, so it knows which lead to pass to {@link
|
||||
* LeadLauncher#relaunch}.
|
||||
* the terminal names no currently-recognised lead. {@link #open} calls this on the calling
|
||||
* lead's own terminal, while it is certainly still live, to resolve {@link
|
||||
* PendingRollover#rolloverKey}. The deferred continuation also calls this, on the OLD terminal,
|
||||
* before tearing it down, so it knows which lead to pass to {@link LeadLauncher#relaunch} — by
|
||||
* that point the live roster may no longer contain the old terminal, so this lookup can return
|
||||
* {@code null} here even though {@link #open}'s earlier call against the same terminal did not.
|
||||
*/
|
||||
private final Function<String, String> leadNameForTerminal;
|
||||
/**
|
||||
@@ -362,14 +373,14 @@ public final class LeadRollover {
|
||||
private final Consumer<Runnable> continuationRunner;
|
||||
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Lead terminal → the token of the roll currently holding that terminal exclusive, for
|
||||
* {@link #confirm}'s single-flight claim. {@link #confirm} claims an entry here with an
|
||||
* atomic put-if-absent once every other gate has passed, refusing with {@link
|
||||
* {@link PendingRollover#rolloverKey} → the token of the roll currently holding that key
|
||||
* exclusive, for {@link #confirm}'s single-flight claim. {@link #confirm} claims an entry here
|
||||
* with an atomic put-if-absent once every other gate has passed, refusing with {@link
|
||||
* RefusalReason#ROLL_ALREADY_RUNNING} when a claim is already held; {@link #runRollover}
|
||||
* releases it in a {@code finally}, on both the success and the thrown-exception path. A
|
||||
* terminal absent from this map has no roll currently in flight for it.
|
||||
* releases it in a {@code finally}, on both the success and the thrown-exception path. A key
|
||||
* absent from this map has no roll currently in flight for it.
|
||||
*/
|
||||
private final Map<String, String> rollingByTerminal = new ConcurrentHashMap<>();
|
||||
private final Map<String, String> rollingByLead = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Finished tokens → what actually happened, for {@link #status}. Bounded by {@link
|
||||
* #OUTCOME_HISTORY_CAP}, oldest evicted first ({@code removeEldestEntry} on an insertion-order
|
||||
@@ -438,8 +449,9 @@ public final class LeadRollover {
|
||||
/**
|
||||
* The lead says it is ready to be replaced. Generates a token and records the resolved
|
||||
* handover path, this moment's wall-clock timestamp (the baseline {@link #confirm} checks the
|
||||
* handover file's modified time against), and {@code leadTerminal} — only that exact terminal
|
||||
* may later {@link #confirm} this token.
|
||||
* handover file's modified time against), {@code leadTerminal} — only that exact terminal may
|
||||
* later {@link #confirm} this token — and {@link PendingRollover#rolloverKey}, resolved here
|
||||
* from {@code leadTerminal} while the calling lead is certainly still live.
|
||||
*
|
||||
* @param leadTerminal the calling lead's terminal id, resolved by the MCP layer from the
|
||||
* connection (see this class's javadoc) — never a client-supplied value
|
||||
@@ -459,7 +471,9 @@ public final class LeadRollover {
|
||||
String token = UUID.randomUUID().toString();
|
||||
long requestedAt = nowMillis.getAsLong();
|
||||
String resolvedPath = resolveHandoverPath(cfg.handoverPath(), leadTerminal);
|
||||
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt);
|
||||
String leadName = leadNameForTerminal.apply(leadTerminal);
|
||||
String rolloverKey = (leadName == null || leadName.isBlank()) ? leadTerminal : leadName;
|
||||
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt, rolloverKey);
|
||||
pending.put(token, p);
|
||||
if (resolvedPath.equals(cfg.handoverPath())) {
|
||||
log.info("lead-rollover: open token={} lead={} handoverPath={} reason={}",
|
||||
@@ -565,12 +579,11 @@ public final class LeadRollover {
|
||||
|
||||
// Single-flight claim: atomic put-if-absent, taken only after every other gate has
|
||||
// passed, so a refused confirm() never takes it. A non-null previous value means a
|
||||
// different, still-running roll already holds this lead terminal.
|
||||
String holder = rollingByTerminal.putIfAbsent(p.leadTerminal(), token);
|
||||
// different, still-running roll already holds this lead's claim.
|
||||
String holder = rollingByLead.putIfAbsent(p.rolloverKey(), token);
|
||||
if (holder != null) {
|
||||
return RollDecision.refused(RefusalReason.ROLL_ALREADY_RUNNING,
|
||||
"lead terminal " + p.leadTerminal() + " already has a roll running under token "
|
||||
+ holder);
|
||||
"lead '" + p.rolloverKey() + "' already has a roll running under token " + holder);
|
||||
}
|
||||
|
||||
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
|
||||
@@ -591,14 +604,14 @@ public final class LeadRollover {
|
||||
} catch (RuntimeException e) {
|
||||
// continuationRunner can reject the hand-off itself (e.g. a bounded executor's
|
||||
// RejectedExecutionException) before runRollover ever starts, so runRollover's own
|
||||
// finally — the only other place that releases rollingByTerminal — never runs either.
|
||||
// finally — the only other place that releases rollingByLead — never runs either.
|
||||
// Release the claim here and overwrite the IN_PROGRESS entry with a terminal outcome,
|
||||
// or this lead terminal could never be rolled again and status() would report
|
||||
// IN_PROGRESS forever for a roll that in fact never started.
|
||||
// or this lead could never be rolled again and status() would report IN_PROGRESS
|
||||
// forever for a roll that in fact never started.
|
||||
log.warn("lead-rollover: continuationRunner rejected token={} lead={}: {} — the roll "
|
||||
+ "never started; releasing its claim and reporting it as FAILED",
|
||||
token, callerTerminal, e.toString(), e);
|
||||
rollingByTerminal.remove(p.leadTerminal(), token);
|
||||
rollingByLead.remove(p.rolloverKey(), token);
|
||||
outcomes.put(token, new RollStatus(RollState.FAILED,
|
||||
"continuationRunner rejected this roll before it ever started: " + e.toString()
|
||||
+ " — the roll never ran; open() a fresh rollover request"));
|
||||
@@ -641,10 +654,10 @@ public final class LeadRollover {
|
||||
+ "a fresh rollover request"));
|
||||
} finally {
|
||||
// Release the single-flight claim on both the normal return and the thrown-exception
|
||||
// path above — a release only on success would leave this lead terminal unrollable
|
||||
// forever after one failure. The conditional two-argument remove only clears the
|
||||
// entry this roll itself holds, never a different roll's claim on the same terminal.
|
||||
rollingByTerminal.remove(p.leadTerminal(), p.token());
|
||||
// path above — a release only on success would leave this lead unrollable forever
|
||||
// after one failure. The conditional two-argument remove only clears the entry this
|
||||
// roll itself holds, never a different roll's claim on the same key.
|
||||
rollingByLead.remove(p.rolloverKey(), p.token());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -5,13 +5,13 @@ import dev.ltms.fleet.herdr.PaneLocator;
|
||||
/**
|
||||
* Resolves <em>who is calling</em> an MCP tool from the connection alone — the anti-spoofing
|
||||
* identity model of the MCP contract. It ties the connection's loopback peer PID (from the OS)
|
||||
* to a herdr agent pane (from herdr), yielding the caller's worker {@code terminal_id}. A caller
|
||||
* that maps to no worker pane — the primary, or an off-host client — resolves to {@code null}.
|
||||
* to a herdr agent pane (from herdr), yielding that pane's {@code terminal_id}. A connection that
|
||||
* maps to no pane resolves to {@code null}; this class assigns no role to either outcome — {@link
|
||||
* dev.ltms.fleet.auth.CallerResolver} does that.
|
||||
*
|
||||
* <p>Both sources are authoritative and unforgeable: the OS reports the real connecting PID, and
|
||||
* herdr owns the PID→pane mapping. A worker cannot claim to be another worker, nor the primary.
|
||||
* Single-host only (the herd shares the {@code fleetd} host); the token path is the split-host
|
||||
* fallback.
|
||||
* herdr owns the PID→pane mapping, so a caller cannot claim to be at another pane. Single-host
|
||||
* only (the herd shares the {@code fleetd} host); the token path is the split-host fallback.
|
||||
*/
|
||||
public final class ConnectionIdentity {
|
||||
|
||||
@@ -46,9 +46,10 @@ public final class ConnectionIdentity {
|
||||
}
|
||||
|
||||
/**
|
||||
* The caller resolved from the connection: its worker {@code terminal} (or {@code null} for the
|
||||
* primary / an off-host client), its {@code pid} (or {@code -1} if not resolvable), and whether
|
||||
* the pane scan behind {@code terminal} ran to completion ({@link #scanComplete}).
|
||||
* The caller resolved from the connection: the {@code terminal} of the pane it connects from
|
||||
* (or {@code null} when the connection maps to no pane), its {@code pid} (or {@code -1} if not
|
||||
* resolvable), and whether the pane scan behind {@code terminal} ran to completion
|
||||
* ({@link #scanComplete}).
|
||||
*/
|
||||
public record Caller(String terminal, long pid, boolean scanComplete) {
|
||||
|
||||
@@ -87,8 +88,8 @@ public final class ConnectionIdentity {
|
||||
}
|
||||
|
||||
/**
|
||||
* The calling worker's {@code terminal_id}, or {@code null} if the caller is not a known
|
||||
* on-host worker (treat as the primary).
|
||||
* The terminal id of the pane the caller connects from, or {@code null} if the connection
|
||||
* maps to no pane.
|
||||
*/
|
||||
public String callerTerminal(String remoteAddr, int remotePort) {
|
||||
return resolve(remoteAddr, remotePort).terminal();
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.Fleetd;
|
||||
import dev.ltms.fleet.auth.AuditLog;
|
||||
import dev.ltms.fleet.auth.Authz;
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
@@ -41,6 +42,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import jakarta.servlet.http.HttpServlet;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Comparator;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -52,6 +54,7 @@ import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.BiFunction;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@@ -116,9 +119,10 @@ public final class FleetMcp {
|
||||
private final ConnectionIdentity identity;
|
||||
/**
|
||||
* Kept as a field (rather than only captured by the {@code contextExtractor} closure) so
|
||||
* {@link #denyFor} can read {@link CallerResolver#knownLeadOrCollaborator()} — the classifier a
|
||||
* collaborator's {@code SEND} is checked against, built from the same lead and collaborator
|
||||
* maps {@link #identity}-based resolution reads.
|
||||
* {@link #denyFor} can read {@link CallerResolver#knownLeadOrCollaborator()} and {@link
|
||||
* CallerResolver#sendableObserverTarget()} — the classifiers a collaborator's and an
|
||||
* observer's {@code SEND} are each checked against, built from the same maps {@link #identity}-
|
||||
* based resolution reads.
|
||||
*/
|
||||
private final CallerResolver callers;
|
||||
private final Metrics metrics; // CB-502: null → auth failures not counted
|
||||
@@ -302,6 +306,29 @@ public final class FleetMcp {
|
||||
public static LeadConfigDirSource none() { return new LeadConfigDirSource(_ -> null, _ -> null); }
|
||||
}
|
||||
|
||||
/**
|
||||
* Pane-discovery facts for {@code fleet_list}'s {@code panes} row — every herdr tab's display
|
||||
* label, the one deliverability gate the status-gated injector itself reads, whether a
|
||||
* terminal is bound to a configured architect slot, and whether an observer caller may
|
||||
* {@code SEND} to it.
|
||||
*
|
||||
* @param tabLabels tab id → its display label, read lazily (only once the row is actually
|
||||
* assembled) since it costs a herdr {@code workspace.list}/{@code tab.list}
|
||||
* scan; a tab herdr reports with no label maps to a {@code null} value
|
||||
* @param deliverable the same gate {@link dev.ltms.fleet.Fleetd#deliverableTo} builds for the
|
||||
* injector, keyed by terminal id — never a second, separately-derived check
|
||||
* @param architectSlot the same classifier {@link CallerResolver#boundToArchitectSlot} resolves
|
||||
* a caller against — never a second, separately-derived check
|
||||
* @param sendableToObserver the same predicate {@link CallerResolver#sendableObserverTarget}
|
||||
* builds for the {@code SEND} gate — never a second, separately-derived
|
||||
* check
|
||||
*/
|
||||
public record PaneSource(Supplier<Map<String, String>> tabLabels, Predicate<String> deliverable,
|
||||
Predicate<String> architectSlot, Predicate<String> sendableToObserver) {
|
||||
/** Inert source — no labels, no deliverable targets, no architect slots, nothing sendable. */
|
||||
public static PaneSource none() { return new PaneSource(Map::of, _ -> false, _ -> false, _ -> false); }
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this
|
||||
* daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the
|
||||
@@ -349,11 +376,9 @@ public final class FleetMcp {
|
||||
* overload instead of failing to compile, and every existing test — none of which exercises
|
||||
* {@code Fleetd.main} itself — stays green while the live daemon quietly answers
|
||||
* {@code NOT_CONFIGURED} to {@code fleet_handover} forever. Collapsing every overload into one
|
||||
* required-everything constructor turns that mistake into a compile error instead: this
|
||||
* project's own antidote for a defaulted parameter surviving as an untested decision (see
|
||||
* {@code FleetdCompletionResolverWiringTest} / {@code FleetdLeadRolloverWiringTest}'s own
|
||||
* javadoc for the same lesson applied to a different seam). A caller that genuinely wants a
|
||||
* feature off must now say so explicitly at the call site — {@code null},
|
||||
* required-everything constructor turns that mistake into a compile error instead, the same
|
||||
* defaulted-parameter guard this codebase applies to other required wiring. A caller that
|
||||
* genuinely wants a feature off must now say so explicitly at the call site — {@code null},
|
||||
* {@link OutageSource#none()}, {@link LeadSeatSource#none()}, {@code List.of()} are all still
|
||||
* perfectly fine values, just never an implicit default reached by omission.
|
||||
*
|
||||
@@ -469,7 +494,10 @@ public final class FleetMcp {
|
||||
recordPrimarySingleton(primaryRegistry, callerTerminal, caller);
|
||||
Map<String, Object> a = req.arguments();
|
||||
String target = str(a, "sessionId");
|
||||
String content = str(a, "content");
|
||||
// An observer's SEND reaches a pane that cannot otherwise distinguish this
|
||||
// from a human paste (see attributeIfObserver); every other caller's content
|
||||
// passes through unchanged.
|
||||
String content = attributeIfObserver(caller, str(a, "content"));
|
||||
String turnId = str(a, "turnId");
|
||||
String coordId = str(a, "coordId");
|
||||
if (coordId != null && !coordId.isBlank()) {
|
||||
@@ -501,8 +529,9 @@ public final class FleetMcp {
|
||||
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), callerOwner);
|
||||
};
|
||||
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
|
||||
// is "is this caller a worker at all", and it can only ever reply as itself.
|
||||
// fleet_reply's identity is the CONNECTION, never an argument. The authz check is
|
||||
// terminal ownership, not a role test: the caller may reply only for its own pane,
|
||||
// which is why no role appears in the check at all.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> replyHandler =
|
||||
(exchange, req) -> {
|
||||
String self = callerTerminal(exchange);
|
||||
@@ -566,13 +595,19 @@ public final class FleetMcp {
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
|
||||
if (denied != null) return denied;
|
||||
// A Supplier: the label lookup costs a herdr scan, and must stay behind
|
||||
// panesVisible so it only runs for a caller that receives the row at all.
|
||||
PaneSource panes = new PaneSource(() -> identity.panes().tabLabelsByTabId(),
|
||||
Fleetd.deliverableTo(presence, callers::leads, callers::collaborators),
|
||||
callers::boundToArchitectSlot, callers.sendableObserverTarget());
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
|
||||
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
callers.collaborators(), collaboratorsVisibleTo(principal(exchange)),
|
||||
new CoordinationSource(leadChannel, peers),
|
||||
coordinatorVisibleTo(principal(exchange)),
|
||||
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)));
|
||||
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)),
|
||||
panes, panesVisibleTo(principal(exchange)), principal(exchange).isObserver());
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
@@ -709,7 +744,8 @@ public final class FleetMcp {
|
||||
if (!authorizationEnforced) {
|
||||
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
|
||||
}
|
||||
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator())) {
|
||||
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator(),
|
||||
callers.sendableObserverTarget())) {
|
||||
if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) {
|
||||
AuditLog.allowed(caller, action, target); // reads would drown the trail
|
||||
}
|
||||
@@ -797,6 +833,17 @@ public final class FleetMcp {
|
||||
return caller.isPrimary() || caller.isArchitect();
|
||||
}
|
||||
|
||||
/**
|
||||
* Who may see {@code fleet_list}'s {@code panes} array — every role that may
|
||||
* {@link Authz.Action#SEND} to some other pane. A plain worker holds {@code READ} but never
|
||||
* {@code SEND}, so it still does not see this array. An observer does hold {@code SEND}, to
|
||||
* another observer pane only, so it sees the array too — but {@code listFleet} filters its rows
|
||||
* to {@link CallerResolver#sendableObserverTarget} and reduces each one; see {@code paneRows}.
|
||||
*/
|
||||
static boolean panesVisibleTo(Principal caller) {
|
||||
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator() || caller.isObserver();
|
||||
}
|
||||
|
||||
/**
|
||||
* The terminal of the caller on this call's connection, or {@code null} when that caller carries
|
||||
* no terminal, which is only the unnamed primary. A named lead, an architect and a worker each
|
||||
@@ -915,6 +962,17 @@ public final class FleetMcp {
|
||||
|
||||
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
|
||||
|
||||
/**
|
||||
* The text an observer's {@code SEND} actually delivers: prefixed with the sender's own
|
||||
* connection-resolved terminal, which the receiving pane cannot otherwise tell apart from a
|
||||
* human paste. Every other caller's content passes through unchanged. Shared with {@code
|
||||
* FleetApp}'s REST entry path so both surfaces attribute identically.
|
||||
*/
|
||||
public static String attributeIfObserver(Principal caller, String content) {
|
||||
return caller != null && caller.isObserver()
|
||||
? "[fleet_send from observer " + caller.terminal() + "]\n" + content : content;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code fleet_send}: delegate {@code content} to a worker session and block for its reply.
|
||||
* The configured profiles are required so a profile name can never bypass target validation.
|
||||
@@ -1968,6 +2026,25 @@ public final class FleetMcp {
|
||||
Map.of(), false, coordination, callerIsPrimary, leadsVisible, membersVisible);
|
||||
}
|
||||
|
||||
/**
|
||||
* As below, with no pane discovery — {@code panes} is {@link PaneSource#none()} and
|
||||
* {@code panesVisible} is {@code false}.
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
LoopHealthSource loopHealth,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
|
||||
LeadConfigDirSource leadConfigDirs,
|
||||
Map<String, String> leads, String selfTerm,
|
||||
Map<String, String> collaborators, boolean collaboratorsVisible,
|
||||
CoordinationSource coordination, boolean callerIsPrimary,
|
||||
boolean leadsVisible, boolean membersVisible) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
|
||||
leadSeats, contextGauge, leadConfigDirs, leads, selfTerm, collaborators, collaboratorsVisible,
|
||||
coordination, callerIsPrimary, leadsVisible, membersVisible, PaneSource.none(), false);
|
||||
}
|
||||
|
||||
/**
|
||||
* The canonical implementation. {@code contextGauge} is the "lead context gauge" (see
|
||||
* {@link LeadContextGauge}) — every wrapper overload above passes a freshly constructed one,
|
||||
@@ -1986,6 +2063,11 @@ public final class FleetMcp {
|
||||
* {@link #leadsVisibleTo})
|
||||
* @param membersVisible whether this caller may see the {@code members} array (see
|
||||
* {@link #membersVisibleTo})
|
||||
* @param panes pane-discovery facts — labels and the deliverable gate for the
|
||||
* {@code panes} row; {@link PaneSource#none()} for a caller that
|
||||
* does not want the row
|
||||
* @param panesVisible whether this caller may see the {@code panes} array (see
|
||||
* {@link #panesVisibleTo})
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
@@ -1996,11 +2078,38 @@ public final class FleetMcp {
|
||||
Map<String, String> leads, String selfTerm,
|
||||
Map<String, String> collaborators, boolean collaboratorsVisible,
|
||||
CoordinationSource coordination, boolean callerIsPrimary,
|
||||
boolean leadsVisible, boolean membersVisible) {
|
||||
boolean leadsVisible, boolean membersVisible,
|
||||
PaneSource panes, boolean panesVisible) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
|
||||
leadSeats, contextGauge, leadConfigDirs, leads, selfTerm, collaborators, collaboratorsVisible,
|
||||
coordination, callerIsPrimary, leadsVisible, membersVisible, panes, panesVisible, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, plus fleetd #758: an observer sees the {@code panes} array too, but filtered to
|
||||
* {@link CallerResolver#sendableObserverTarget} and each row reduced to the five fields an
|
||||
* observer may learn — see {@code paneRows}/{@code paneRow}.
|
||||
*
|
||||
* @param callerIsObserver whether the {@code fleet_list} caller is an observer; every wrapper
|
||||
* overload above passes {@code false}, so a test that wants the
|
||||
* filtered, reduced view must call this overload with an explicit
|
||||
* {@code true}
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
LoopHealthSource loopHealth,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
|
||||
LeadConfigDirSource leadConfigDirs,
|
||||
Map<String, String> leads, String selfTerm,
|
||||
Map<String, String> collaborators, boolean collaboratorsVisible,
|
||||
CoordinationSource coordination, boolean callerIsPrimary,
|
||||
boolean leadsVisible, boolean membersVisible,
|
||||
PaneSource panes, boolean panesVisible, boolean callerIsObserver) {
|
||||
try {
|
||||
// Neither row's assembly (leadView/memberCapacityView probing herdr for live status)
|
||||
// runs unless at least one of them needs the live-agent lookup backing it.
|
||||
Map<String, Agent> live = (leadsVisible || membersVisible)
|
||||
Map<String, Agent> live = (leadsVisible || membersVisible || panesVisible)
|
||||
? workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
@@ -2044,6 +2153,10 @@ public final class FleetMcp {
|
||||
.map(e -> collaboratorRow(e.getKey(), e.getValue()))
|
||||
.toList());
|
||||
}
|
||||
// gate BEFORE assembling the row, so the key is absent rather than present-and-empty.
|
||||
if (panesVisible) {
|
||||
result.put("panes", paneRows(live, roster, leads, collaborators, panes, callerIsObserver));
|
||||
}
|
||||
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
|
||||
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
|
||||
// key is absent rather than present-and-empty.
|
||||
@@ -2369,6 +2482,102 @@ public final class FleetMcp {
|
||||
return m;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code panes.tabLabels()}'s herdr scan, or an empty map on a {@code HerdrException} — a
|
||||
* missing label must not cost the {@code leads}/{@code members}/{@code capacity}/
|
||||
* {@code coordinator} rows that share {@code listFleet}'s own {@code catch}.
|
||||
*/
|
||||
private static Map<String, String> tabLabelsOrEmpty(PaneSource panes) {
|
||||
try {
|
||||
return panes.tabLabels().get();
|
||||
} catch (HerdrException e) {
|
||||
return Map.of();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* One row per herdr-tracked agent pane, sorted by terminal id for a stable order. {@code live}
|
||||
* is the same terminal-keyed {@link Agent} map {@code leadView}/{@code memberCapacityView}
|
||||
* already read, so a pane neither configured as a lead nor spawned as a member — a hand-opened
|
||||
* tab — still gets a row here.
|
||||
*
|
||||
* <p>fleetd #758: for an observer caller ({@code observerView}), the rows are filtered to
|
||||
* {@link PaneSource#sendableToObserver} before being built, and each row is reduced — see
|
||||
* {@code paneRow}.
|
||||
*/
|
||||
private static List<Map<String, Object>> paneRows(Map<String, Agent> live, List<MemberSession> roster,
|
||||
Map<String, String> leads, Map<String, String> collaborators, PaneSource panes,
|
||||
boolean observerView) {
|
||||
final Map<String, String> tabLabels = tabLabelsOrEmpty(panes);
|
||||
Map<String, MemberSession> byTerminal = roster.stream()
|
||||
.filter(s -> s.terminalId() != null)
|
||||
.collect(Collectors.toMap(MemberSession::terminalId, Function.identity(), (_, b) -> b));
|
||||
return live.values().stream()
|
||||
.filter(a -> !observerView || panes.sendableToObserver().test(a.terminalId()))
|
||||
.sorted(Comparator.comparing(Agent::terminalId))
|
||||
.map(a -> paneRow(a, byTerminal.get(a.terminalId()), leads, collaborators, tabLabels, panes,
|
||||
observerView))
|
||||
.toList();
|
||||
}
|
||||
|
||||
/**
|
||||
* @param session the roster entry for this pane's terminal, or {@code null} for a pane the
|
||||
* daemon never spawned as a member (a hand-opened tab, or a configured lead)
|
||||
* @param tabLabels tab id → its herdr display label; a tab absent here, or carrying a
|
||||
* {@code null} label itself, projects as a {@code null} "label"
|
||||
* @param observerView fleetd #758: an observer's row carries only {@code sessionId}, {@code
|
||||
* label}, {@code status}, {@code role}, {@code deliverable} — never {@code
|
||||
* paneId} (the {@code fleet_stop} handle), {@code workspaceId}, {@code tabId},
|
||||
* {@code agentType}, or {@code cwd} (a member's worktree path is the lead's
|
||||
* business)
|
||||
*/
|
||||
private static Map<String, Object> paneRow(Agent a, MemberSession session, Map<String, String> leads,
|
||||
Map<String, String> collaborators, Map<String, String> tabLabels, PaneSource panes,
|
||||
boolean observerView) {
|
||||
Map<String, Object> m = new LinkedHashMap<>();
|
||||
m.put("sessionId", a.terminalId());
|
||||
if (!observerView) {
|
||||
m.put("paneId", a.paneId());
|
||||
m.put("workspaceId", a.workspaceId());
|
||||
m.put("tabId", a.tabId());
|
||||
}
|
||||
m.put("label", a.tabId() == null ? null : tabLabels.get(a.tabId()));
|
||||
if (!observerView) {
|
||||
m.put("agentType", a.agentType());
|
||||
}
|
||||
m.put("status", a.status() == null ? "unknown" : a.status().name().toLowerCase());
|
||||
m.put("role", paneRole(a.terminalId(), session, leads, collaborators, panes));
|
||||
m.put("deliverable", panes.deliverable().test(a.terminalId()));
|
||||
if (!observerView && session != null && session.cwd() != null) {
|
||||
m.put("cwd", session.cwd());
|
||||
}
|
||||
return m;
|
||||
}
|
||||
|
||||
/**
|
||||
* The role this pane resolves as: a spawned member's own {@link MemberRole}, else "lead" for a
|
||||
* configured but currently-unoccupied lead pane, else "architect" for a pane bound to a
|
||||
* configured architect slot with no live member session, else "collaborator" for a configured
|
||||
* but currently-unoccupied collaborator tab, else "observer" for a pane this daemon neither
|
||||
* spawned nor configured.
|
||||
*/
|
||||
private static String paneRole(String terminal, MemberSession session, Map<String, String> leads,
|
||||
Map<String, String> collaborators, PaneSource panes) {
|
||||
if (session != null) {
|
||||
return session.role().wireName();
|
||||
}
|
||||
if (leads.containsKey(terminal)) {
|
||||
return "lead";
|
||||
}
|
||||
if (panes.architectSlot().test(terminal)) {
|
||||
return "architect";
|
||||
}
|
||||
if (collaborators.containsKey(terminal)) {
|
||||
return "collaborator";
|
||||
}
|
||||
return "observer";
|
||||
}
|
||||
|
||||
/**
|
||||
* Reads the lead context gauge for one lead. {@code configDir} is this lead's configured
|
||||
* {@code CLAUDE_CONFIG_DIR} override (see {@link LeadConfigDirSource}), derived from
|
||||
@@ -2574,8 +2783,8 @@ public final class FleetMcp {
|
||||
+ "orchestrators, each with its sessionId (the address to fleet_send to), "
|
||||
+ "name, live status, and 'self': true on your own row; this is how you "
|
||||
+ "discover a peer lead without being told its address. 'members' are the "
|
||||
+ "sessions delegated to — each with sessionId, paneId, role (architect/dev/"
|
||||
+ "reviewer), profile (the backend it runs on), state, optional "
|
||||
+ "sessions delegated to — each with sessionId, paneId, role (" + MemberRole.wireNames()
|
||||
+ "), profile (the backend it runs on), state, optional "
|
||||
+ "worktree/branch/owner/agentSessionId, and live herdr status. agentSessionId, "
|
||||
+ "when present, is the id to pass as fleet_spawn's resumeSessionId to relaunch "
|
||||
+ "onto that same conversation. It is ABSENT — not a guess — for a member fleetd "
|
||||
|
||||
@@ -445,27 +445,6 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
+ " distinct candidate(s): " + String.join(", ", unreachable));
|
||||
}
|
||||
|
||||
/**
|
||||
* Refuse an explicit-profile spawn when the profile is at its {@code maxLoad} cap.
|
||||
*
|
||||
* <p>maxLoad is a documented, unconditional capacity limit (see {@code FleetConfig.Profile#maxLoad}),
|
||||
* and the charter makes explicit-profile spawns the normal path — so enforcing it only in placement
|
||||
* ({@code PlacementPolicyUtil}, package-private, hence not linked) would leave the cap dead config
|
||||
* on every call that names a profile. Same rule as placement: {@code live >= cap} is at capacity.
|
||||
*
|
||||
* <p>Deliberately no fallback to another profile: the caller named {@code profile} for a cost/model
|
||||
* reason, and silently re-routing a paid-tier (subscription) request elsewhere is worse than
|
||||
* refusing it. A caller that wants placement should omit the profile and let the policy pick.
|
||||
*
|
||||
* <p>Known TOCTOU limitation — documented, not fixed. {@link #liveCount} is read outside any lock and
|
||||
* {@code SessionManager} registers a session only after {@code launcher.spawn} returns, so two
|
||||
* genuinely concurrent spawns can both pass this check. The race already exists on the placement
|
||||
* path. Closing it needs slot reservation in the registry; serializing spawn here would block on
|
||||
* the readiness gate and is a far worse trade.
|
||||
*
|
||||
* @param profile the profile the caller explicitly named
|
||||
* @throws PlacementException when the profile is at capacity
|
||||
*/
|
||||
/**
|
||||
* Refuse an explicit-profile spawn whose credential is quarantined (CB-578 stage B): a prior
|
||||
* {@code BACKEND_EXHAUSTED} classification on this profile, or on another profile sharing its
|
||||
@@ -527,6 +506,27 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
.collect(Collectors.toSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* Refuse an explicit-profile spawn when the profile is at its {@code maxLoad} cap.
|
||||
*
|
||||
* <p>maxLoad is a documented, unconditional capacity limit (see {@code FleetConfig.Profile#maxLoad}),
|
||||
* and the charter makes explicit-profile spawns the normal path — so enforcing it only in placement
|
||||
* ({@code PlacementPolicyUtil}, package-private, hence not linked) would leave the cap dead config
|
||||
* on every call that names a profile. Same rule as placement: {@code live >= cap} is at capacity.
|
||||
*
|
||||
* <p>No fallback to another profile: the caller named {@code profile} for a cost/model
|
||||
* reason, and silently re-routing a paid-tier (subscription) request elsewhere is worse than
|
||||
* refusing it. A caller that wants placement should omit the profile and let the policy pick.
|
||||
*
|
||||
* <p>Known TOCTOU limitation — documented, not fixed. {@link #liveCount} is read outside any lock and
|
||||
* {@code SessionManager} registers a session only after {@code launcher.spawn} returns, so two
|
||||
* genuinely concurrent spawns can both pass this check. The race already exists on the placement
|
||||
* path. Closing it needs slot reservation in the registry; serializing spawn here would block on
|
||||
* the readiness gate and is a far worse trade.
|
||||
*
|
||||
* @param profile the profile the caller explicitly named
|
||||
* @throws PlacementException when the profile is at capacity
|
||||
*/
|
||||
private void enforceMaxLoad(String profile) {
|
||||
// Absent config, or a config whose maxLoad normalized to null (ABSENT ⇒ unlimited at load),
|
||||
// means no cap — never cap what wasn't configured. Note "non-positive ⇒ unlimited" was true
|
||||
|
||||
@@ -474,7 +474,6 @@ public final class EnvAllowListScrub {
|
||||
}
|
||||
}
|
||||
|
||||
/** Best-effort recursive delete; failures are swallowed — JVM-exit cleanup is the backstop. */
|
||||
/**
|
||||
* Remove generated directories left behind by an earlier daemon process.
|
||||
*
|
||||
@@ -511,6 +510,7 @@ public final class EnvAllowListScrub {
|
||||
}
|
||||
}
|
||||
|
||||
/** Best-effort recursive delete; failures are swallowed — JVM-exit cleanup is the backstop. */
|
||||
static void deleteRecursively(Path dir) {
|
||||
if (dir == null || !Files.exists(dir)) {
|
||||
return;
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.PromptBox;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -23,7 +24,8 @@ import java.util.function.Supplier;
|
||||
* <p><strong>Status-gated, exactly like {@link ReplyPushLoop}.</strong> A pane may only be injected
|
||||
* into at a turn boundary ({@link AgentStatus#injectable()} — idle, blocked or done); pasting into
|
||||
* a live turn corrupts it. So a tick that finds the lead busy simply does nothing and comes back
|
||||
* later.
|
||||
* later. The same holds for a lead whose prompt box holds unsubmitted text ({@link PromptBox}) —
|
||||
* delivering there would submit the operator's half-typed line along with the message.
|
||||
*
|
||||
* <p><strong>Ack only after delivery.</strong> A message is acked — removed from the broker — only
|
||||
* once {@link AgentControl#send} has actually put it in the pane. Anything not delivered (no lead
|
||||
@@ -52,6 +54,7 @@ public final class LeadCoordLoop {
|
||||
|
||||
private final LeadChannel channel;
|
||||
private final AgentControl agents;
|
||||
private final PromptBox promptBox;
|
||||
private final Supplier<Map<String, String>> leads;
|
||||
private final ScheduledExecutorService scheduler;
|
||||
private final long intervalMs;
|
||||
@@ -73,6 +76,7 @@ public final class LeadCoordLoop {
|
||||
ScheduledExecutorService scheduler, long intervalMs) {
|
||||
this.channel = channel;
|
||||
this.agents = agents;
|
||||
this.promptBox = new PromptBox(agents);
|
||||
this.leads = leads;
|
||||
this.scheduler = scheduler;
|
||||
this.intervalMs = intervalMs;
|
||||
@@ -149,6 +153,11 @@ public final class LeadCoordLoop {
|
||||
lead, status, held.size());
|
||||
return;
|
||||
}
|
||||
if (!promptBox.clearToSubmit(lead)) {
|
||||
log.debug("lead coordination: lead {} has unsubmitted text in its prompt box, holding {} message(s)",
|
||||
lead, held.size());
|
||||
return;
|
||||
}
|
||||
try {
|
||||
agents.send(lead, DELIVERY_FORMAT.formatted(msg.from(), msg.content()));
|
||||
} catch (RuntimeException e) {
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.PromptBox;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.metrics.FleetMetrics;
|
||||
@@ -33,7 +34,7 @@ import java.util.function.Supplier;
|
||||
* no such block it is never constructed, so upgrading the daemon cannot silently acquire a behaviour
|
||||
* that spends the operator's model subscription on its own initiative (constraint 1).
|
||||
*
|
||||
* <p>Four invariants keep it from becoming a runaway subscription burner:
|
||||
* <p>Five invariants keep it from becoming a runaway subscription burner:
|
||||
* <ol>
|
||||
* <li><b>Status-gated</b> — a {@code WORKING} lead is making progress and is never touched; only an
|
||||
* injectable (idle/done/blocked) lead is even considered (constraint 2).</li>
|
||||
@@ -45,6 +46,8 @@ import java.util.function.Supplier;
|
||||
* <li><b>Never races {@link ReplyPushLoop}</b> — while that loop is actively nudging any target this
|
||||
* loop stands down, so two competing injections never start two turns in the same pane
|
||||
* (constraint 6).</li>
|
||||
* <li><b>Never submits the operator's draft</b> — a nudge is held while the lead's prompt box holds
|
||||
* unsubmitted text ({@link PromptBox}), because the delivery pastes and submits in one call.</li>
|
||||
* </ol>
|
||||
*
|
||||
* <p><b>fleetd #609 — context-high notice.</b> Optionally ({@code contextHighNudge}, opt-in like the
|
||||
@@ -65,6 +68,7 @@ public final class LeadHeartbeatLoop {
|
||||
|
||||
private final PrimaryRegistry primaryRegistry;
|
||||
private final AgentControl agents;
|
||||
private final PromptBox promptBox;
|
||||
private final ReplyInbox inbox;
|
||||
private final Supplier<List<MemberSession>> roster;
|
||||
private final ReplyPushLoop pushLoop;
|
||||
@@ -135,6 +139,7 @@ public final class LeadHeartbeatLoop {
|
||||
boolean requireOperatorConfirm) {
|
||||
this.primaryRegistry = primaryRegistry;
|
||||
this.agents = agents;
|
||||
this.promptBox = new PromptBox(agents);
|
||||
this.inbox = inbox;
|
||||
this.roster = roster;
|
||||
this.pushLoop = pushLoop;
|
||||
@@ -187,7 +192,9 @@ public final class LeadHeartbeatLoop {
|
||||
/** Idle past the quiet period with nothing pending and the cap exhausted — stop until new state appears. */
|
||||
QUIET_DONE,
|
||||
/** {@link ReplyPushLoop} is actively nudging — stand aside rather than start a competing turn. */
|
||||
STAND_DOWN
|
||||
STAND_DOWN,
|
||||
/** The lead's prompt box holds unsubmitted text — hold the nudge rather than submit that text. */
|
||||
DRAFT_HELD
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -330,6 +337,7 @@ public final class LeadHeartbeatLoop {
|
||||
idleSinceNanos == NOT_IDLE ? null : idleSinceNanos,
|
||||
quietCount, status, pushLoop.isActive(), leadKnown, fleet,
|
||||
reading.state(), contextNotified);
|
||||
d = holdIfOperatorIsTyping(d);
|
||||
applyDecision(d);
|
||||
switch (d.action()) {
|
||||
case INJECT -> injectNudge(d, fleet, reading);
|
||||
@@ -337,11 +345,33 @@ public final class LeadHeartbeatLoop {
|
||||
countNudge("exhausted");
|
||||
contextNotified = d.contextNotified();
|
||||
}
|
||||
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> contextNotified = d.contextNotified();
|
||||
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN, DRAFT_HELD -> contextNotified = d.contextNotified();
|
||||
}
|
||||
scheduleNext();
|
||||
}
|
||||
|
||||
/**
|
||||
* Turn a decision to inject into {@link Action#DRAFT_HELD} when the lead's prompt box holds text
|
||||
* the operator has not submitted. The pane read happens only for a decision that would otherwise
|
||||
* send, so a busy or debouncing lead costs no extra herdr call.
|
||||
*
|
||||
* <p>The held decision carries this tick's idle window but the <em>pre-tick</em> quiet count and
|
||||
* context latch: nothing reached the pane, so neither the quiet budget nor the one context notice
|
||||
* per stretch may be spent on it.
|
||||
*/
|
||||
private Decision holdIfOperatorIsTyping(Decision d) {
|
||||
if (d.action() != Action.INJECT) {
|
||||
return d;
|
||||
}
|
||||
var lead = primaryRegistry.currentPrimaryTerminal();
|
||||
if (lead.isEmpty() || promptBox.clearToSubmit(lead.get())) {
|
||||
return d;
|
||||
}
|
||||
log.debug("idle-heartbeat: lead {} has unsubmitted text in its prompt box, holding the nudge",
|
||||
lead.get());
|
||||
return new Decision(Action.DRAFT_HELD, d.idleSinceNanos(), quietCount, contextNotified);
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist the idle/quiet state a decision returned, so the next tick starts from it.
|
||||
*
|
||||
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.msg;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.PromptBox;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.metrics.FleetMetrics;
|
||||
import dev.ltms.fleet.metrics.Metrics;
|
||||
@@ -50,6 +51,11 @@ import java.util.stream.Collectors;
|
||||
* exhausting its cap does not stop nudges about the others (post-CB-590 regression fix; see
|
||||
* {@link #decide}) — whichever the durable inbox / pending set doesn't already answer via
|
||||
* {@code STOP}.
|
||||
*
|
||||
* <p>A lead that is injectable is nudged only when its prompt box is also empty
|
||||
* ({@link PromptBox}): the delivery pastes and submits in one call, so a nudge into a box holding
|
||||
* the operator's half-typed line would submit that line too. A nudge held for that reason waits for
|
||||
* the next tick like any other, and the pending work is re-read then.
|
||||
*/
|
||||
public final class ReplyPushLoop {
|
||||
|
||||
@@ -81,6 +87,7 @@ public final class ReplyPushLoop {
|
||||
|
||||
private final PrimaryRegistry primaryRegistry;
|
||||
private final AgentControl agents;
|
||||
private final PromptBox promptBox;
|
||||
private final ReplyInbox inbox;
|
||||
private final ScheduledExecutorService scheduler;
|
||||
private final int maxReminders;
|
||||
@@ -125,6 +132,7 @@ public final class ReplyPushLoop {
|
||||
int maxReminders, long backoffMs, Metrics metrics) {
|
||||
this.primaryRegistry = primaryRegistry;
|
||||
this.agents = agents;
|
||||
this.promptBox = new PromptBox(agents);
|
||||
this.inbox = inbox;
|
||||
this.scheduler = scheduler;
|
||||
this.maxReminders = maxReminders;
|
||||
@@ -398,11 +406,15 @@ public final class ReplyPushLoop {
|
||||
log.debug("push: status check failed for lead {}, will retry", lead, e);
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
if (status.injectable()) {
|
||||
return Action.INJECT;
|
||||
if (!status.injectable()) {
|
||||
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
|
||||
return Action.WAIT_BUSY;
|
||||
if (!promptBox.clearToSubmit(lead)) {
|
||||
log.debug("push: lead {} has unsubmitted text in its prompt box, waiting", lead);
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
return Action.INJECT;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package dev.ltms.fleet.peer;
|
||||
|
||||
import java.util.Locale;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
/**
|
||||
* What a member is <em>for</em> — the contract it runs under.
|
||||
@@ -100,6 +102,11 @@ public enum MemberRole {
|
||||
return null;
|
||||
}
|
||||
|
||||
/** The wire name of every role, joined with {@code ", "} in declaration order. */
|
||||
public static String wireNames() {
|
||||
return Stream.of(values()).map(MemberRole::wireName).collect(Collectors.joining(", "));
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a config/wire spelling, case-insensitively.
|
||||
*
|
||||
@@ -118,14 +125,7 @@ public enum MemberRole {
|
||||
}
|
||||
}
|
||||
}
|
||||
StringBuilder valid = new StringBuilder();
|
||||
for (MemberRole r : values()) {
|
||||
if (!valid.isEmpty()) {
|
||||
valid.append(", ");
|
||||
}
|
||||
valid.append(r.wireName());
|
||||
}
|
||||
throw new IllegalArgumentException(
|
||||
"unknown member role '" + s + "'; valid roles are: " + valid);
|
||||
"unknown member role '" + s + "'; valid roles are: " + wireNames());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.inject.MemberPresence;
|
||||
import dev.ltms.fleet.member.MemberCredentialPolicyView;
|
||||
import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
@@ -281,6 +282,18 @@ public final class FleetApp {
|
||||
return Authz.permits(caller, action, target, knownLeadOrCollaborator);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #permitsFor(Principal, Authz.Action, String, Predicate)}, also threading the
|
||||
* classifier an observer's {@code SEND} is checked against; pass {@link #auth}'s own
|
||||
* {@code sendableObserverTarget()} to exercise the real production gate, as {@link #allow}
|
||||
* does.
|
||||
*/
|
||||
static boolean permitsFor(Principal caller, Authz.Action action, String target,
|
||||
Predicate<String> knownLeadOrCollaborator,
|
||||
Predicate<String> knownObserverTarget) {
|
||||
return Authz.permits(caller, action, target, knownLeadOrCollaborator, knownObserverTarget);
|
||||
}
|
||||
|
||||
/**
|
||||
* Gate a handler on the CB-505 authorization table. Returns {@code true} when the request may
|
||||
* proceed; otherwise writes the error response and returns {@code false}.
|
||||
@@ -294,7 +307,7 @@ public final class FleetApp {
|
||||
return true; // legacy: authorization not enforced
|
||||
}
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator())) {
|
||||
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.sendableObserverTarget())) {
|
||||
if (action != Authz.Action.READ && action != Authz.Action.METRICS
|
||||
&& action != Authz.Action.TASK_READ) {
|
||||
AuditLog.allowed(caller, action, target); // reads would drown the trail
|
||||
@@ -436,14 +449,27 @@ public final class FleetApp {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Tab id → its herdr display label, or an empty map on a {@code workspace.list}/{@code
|
||||
* tab.list} failure — a missing label must not cost the agent roster.
|
||||
*/
|
||||
private Map<String, String> tabLabelsOrEmpty() {
|
||||
try {
|
||||
return new PaneLocator(herdr, memberHerdr).tabLabelsByTabId();
|
||||
} catch (HerdrException e) {
|
||||
return Map.of();
|
||||
}
|
||||
}
|
||||
|
||||
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
|
||||
private void agents(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /agents"), null)) {
|
||||
return;
|
||||
}
|
||||
final Map<String, String> tabLabels = tabLabelsOrEmpty();
|
||||
try {
|
||||
ctx.status(200).json(Map.of("agents",
|
||||
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
|
||||
workers.list().stream().map(Agent.class::cast).map(a -> view(a, tabLabels)).toList()));
|
||||
} catch (HerdrException e) {
|
||||
// fleetd #297: workers.list() reaches herdr — a transport failure must land in the same
|
||||
// {error, detail} envelope every other failure path here uses, not escape as a bare
|
||||
@@ -686,6 +712,10 @@ public final class FleetApp {
|
||||
ctx.status(400).json(Map.of("error", "bad_request", "detail", "content is required"));
|
||||
return;
|
||||
}
|
||||
// An observer's SEND reaches a pane that cannot otherwise distinguish this from a human
|
||||
// paste (see FleetMcp#attributeIfObserver, the same rule on the MCP entry path); every
|
||||
// other caller's content passes through unchanged.
|
||||
content = FleetMcp.attributeIfObserver(caller, content);
|
||||
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
|
||||
|
||||
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
|
||||
@@ -935,13 +965,19 @@ public final class FleetApp {
|
||||
}
|
||||
}
|
||||
|
||||
/** Stable JSON projection of an agent (null-safe for the start-time shape). */
|
||||
private static Map<String, Object> view(Agent a) {
|
||||
/**
|
||||
* Stable JSON projection of an agent (null-safe for the start-time shape).
|
||||
*
|
||||
* @param tabLabels tab id → its herdr display label; a tab absent from this map, or carrying
|
||||
* a {@code null} label itself, projects as a {@code null} "label"
|
||||
*/
|
||||
private static Map<String, Object> view(Agent a, Map<String, String> tabLabels) {
|
||||
Map<String, Object> m = new LinkedHashMap<>();
|
||||
m.put("terminalId", a.terminalId());
|
||||
m.put("paneId", a.paneId());
|
||||
m.put("workspaceId", a.workspaceId());
|
||||
m.put("tabId", a.tabId());
|
||||
m.put("label", tabLabels.get(a.tabId()));
|
||||
m.put("sessionId", a.sessionId());
|
||||
m.put("agentType", a.agentType());
|
||||
m.put("status", a.status().name().toLowerCase());
|
||||
|
||||
@@ -442,38 +442,6 @@ public final class GitWorktrees implements Worktrees {
|
||||
ENVIRONMENT_CREDENTIAL_HELPER);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #configureEnvironmentCredentialHelper} only ever fires for an HTTPS origin — Git never
|
||||
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
|
||||
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member sitting on that origin never
|
||||
* reaches the helper and the repo-scoped {@code WORKER_GITEA_TOKEN} is simply not used.
|
||||
*
|
||||
* <p>An earlier version of this javadoc justified the rewrite by claiming a member <em>cannot</em>
|
||||
* push once {@code memberCredentials.policy: allow-list} blocks {@code SSH_AUTH_SOCK}, because
|
||||
* "there is no private key file on this host, only an ssh-agent socket". That premise is false
|
||||
* (fleetd #184): {@code ssh -G} resolves a readable, passphrase-free {@code IdentityFile} outside
|
||||
* {@code ~/.ssh}, and a member — same OS user — pushes over SSH with the socket blanked. The
|
||||
* rewrite is still worth having, but for the reason below rather than that one: it routes the
|
||||
* member through its own scoped token instead of the operator's ssh identity, which is what makes
|
||||
* a member's pushes attributable and revocable.
|
||||
*
|
||||
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
|
||||
* <ssh-base>}, set with {@code --worktree} so it lands only in
|
||||
* {@code <worktree>/.git/worktrees/<name>/config.worktree} (enabled by
|
||||
* {@code extensions.worktreeConfig}, already turned on above) and never touches the shared
|
||||
* repo-level config the primary checkout also reads. {@code insteadOf} — not
|
||||
* {@code pushInsteadOf} — because a member may also need to fetch or rebase, and both should go
|
||||
* through the member's own token for the same reason.
|
||||
*
|
||||
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
|
||||
* itself — never a hardcoded forge host, which is exactly what #177 removed. An origin that is
|
||||
* already {@code https://} is left alone; the credential helper already covers it. An origin
|
||||
* that is neither {@code ssh://} nor {@code https://} — including the scp-like shorthand
|
||||
* ({@code git@host:path}, no scheme) — is left untouched deliberately: that shorthand's
|
||||
* {@code host:path} split is defined by the user's ssh_config aliases, not by URI syntax, so
|
||||
* guessing at it risks rewriting to the wrong place. A repo provisioned from that form keeps
|
||||
* today's (broken, if the policy blocks the agent) SSH-only behaviour rather than a wrong rewrite.
|
||||
*/
|
||||
/**
|
||||
* Blank the user-info of a remote URL before it reaches a log. A remote URL is not obviously a
|
||||
* credential channel, which is exactly why one has leaked here three times ({@code git remote -v}
|
||||
@@ -485,6 +453,30 @@ public final class GitWorktrees implements Worktrees {
|
||||
return url == null ? null : url.replaceAll("://[^@/]*@", "://<redacted>@");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #configureEnvironmentCredentialHelper} only fires for an HTTPS origin — Git never
|
||||
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
|
||||
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member on that origin never
|
||||
* reaches the helper, and the repo-scoped token goes unused without a separate rewrite.
|
||||
*
|
||||
* <p>This method routes the member through its own scoped token instead of the operator's ssh
|
||||
* identity, which is what makes a member's pushes attributable and revocable.
|
||||
*
|
||||
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
|
||||
* <ssh-base>}, set with {@code --worktree} so it lands only in
|
||||
* {@code <worktree>/.git/worktrees/<name>/config.worktree} and never touches the shared
|
||||
* repo-level config the primary checkout also reads. {@code insteadOf} — not
|
||||
* {@code pushInsteadOf} — because a member may also need to fetch or rebase through its own
|
||||
* token.
|
||||
*
|
||||
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
|
||||
* itself, never a hardcoded forge host. An origin already {@code https://} is left alone; the
|
||||
* credential helper already covers it. An origin that is neither {@code ssh://} nor
|
||||
* {@code https://} — including the scp-like shorthand ({@code git@host:path}, no scheme) — is
|
||||
* left untouched: that shorthand's {@code host:path} split is defined by the user's ssh_config
|
||||
* aliases, not by URI syntax, so guessing at it risks rewriting to the wrong place, and that
|
||||
* origin keeps SSH-only push behaviour instead.
|
||||
*/
|
||||
private void configureHttpsUrlRewriteForSshOrigin(String repoRoot, String worktreePath) {
|
||||
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
|
||||
return;
|
||||
|
||||
@@ -91,6 +91,11 @@ class FleetdBackendErrorSinkTest {
|
||||
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", "term_primary").put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
// The lead-nudge paths read the input box before pasting into it.
|
||||
return MAPPER.createObjectNode().set("read",
|
||||
MAPPER.createObjectNode().put("text", FakeHerdr.IDLE_PROMPT_CARET));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
prompts.add(params);
|
||||
sendLatch.countDown();
|
||||
|
||||
@@ -1,96 +1,174 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import com.tngtech.archunit.base.DescribedPredicate;
|
||||
import com.tngtech.archunit.core.domain.Dependency;
|
||||
import com.tngtech.archunit.core.domain.JavaClass;
|
||||
import com.tngtech.archunit.core.domain.JavaClass.Predicates;
|
||||
import com.tngtech.archunit.core.domain.JavaClasses;
|
||||
import com.tngtech.archunit.core.importer.ClassFileImporter;
|
||||
import com.tngtech.archunit.core.importer.ImportOption;
|
||||
import com.tngtech.archunit.library.dependencies.SliceRule;
|
||||
import com.tngtech.archunit.library.dependencies.SlicesRuleDefinition;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.fail;
|
||||
|
||||
/**
|
||||
* fleetd #131 (CB-627): enforce package boundaries with an ArchUnit test instead of a
|
||||
* Maven module split.
|
||||
* Enforces package boundaries between the top-level {@code dev.ltms.fleet.*} packages.
|
||||
*
|
||||
* <p>This test fails the build the moment a NEW cycle appears between the top-level
|
||||
* {@code dev.ltms.fleet.*} packages. Today's cycles are recorded below as explicit,
|
||||
* narrow exceptions: each one ignores dependencies between exactly the two named
|
||||
* packages, in both directions, and nothing else. A cycle through any other pair of
|
||||
* packages -- or a brand new pair -- still fails this test.
|
||||
* <p>{@link #BASELINE_EDGES} names the exact {@code origin class -> target class}
|
||||
* dependencies allowed to cross a top-level package boundary. Any dependency between two
|
||||
* top-level packages that is not in that set fails this test, including a brand new
|
||||
* dependency between a pair of packages that already has other baselined edges. A baseline
|
||||
* entry whose dependency no longer exists in the code also fails this test, so the baseline
|
||||
* always names exactly today's exceptions and nothing more.
|
||||
*
|
||||
* <p><b>Main code only.</b> The import excludes test classes
|
||||
* ({@link ImportOption.Predefined#DO_NOT_INCLUDE_TESTS}). Test code legitimately wires
|
||||
* across many packages for setup and mocking; that is not part of the shipped
|
||||
* architecture this rule protects. Verified: importing test classes too pulls in a much
|
||||
* larger, noisier cycle set -- {@code herdr}, {@code member}, {@code peer}, {@code
|
||||
* config}, {@code guard} and {@code placement} all show up in cycles that disappear the
|
||||
* moment test classes are excluded. Scanning off the classpath via {@code
|
||||
* importPackages(...)} (not a hardcoded {@code target/classes} path) also keeps this
|
||||
* test correct regardless of the working directory the build is invoked from.
|
||||
*
|
||||
* <p><b>No package moves here</b> -- ticket #131 is explicit that removing a cycle is
|
||||
* its own, later PR. See the comment on each exception below for which ticket step
|
||||
* removes it.
|
||||
* ({@link ImportOption.Predefined#DO_NOT_INCLUDE_TESTS}). Scanning off the classpath via
|
||||
* {@code importPackages(...)} keeps this test correct regardless of the working directory
|
||||
* the build is invoked from.
|
||||
*/
|
||||
class PackageCyclesTest {
|
||||
|
||||
/**
|
||||
* Exact {@code "origin -> target"} class dependencies allowed to cross a top-level
|
||||
* package boundary. Each entry is one directed edge between two specific classes; a
|
||||
* two-way relationship between a pair of packages is listed as two separate entries,
|
||||
* one per direction.
|
||||
*/
|
||||
private static final Set<String> BASELINE_EDGES = Set.of(
|
||||
"dev.ltms.fleet.auth.CallerResolver -> dev.ltms.fleet.mcp.ConnectionIdentity",
|
||||
"dev.ltms.fleet.auth.CallerResolver -> dev.ltms.fleet.mcp.ConnectionIdentity$Caller",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.AuditLog",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Authz",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Authz$Action",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.CallerResolver",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Principal",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Role",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadChannel",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadChannel$MailboxState",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadMessage",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$AskOutcome",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$AskResult",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Outcome",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Outstanding",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$PendingAsk",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Phase",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Reply",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$ReplyOutcome",
|
||||
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$TaskView",
|
||||
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$AskOutcome",
|
||||
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$Outcome",
|
||||
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$Phase",
|
||||
"dev.ltms.fleet.mcp.FleetMcp$CoordinationSource -> dev.ltms.fleet.msg.LeadChannel",
|
||||
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.mcp.PrimaryRegistry",
|
||||
"dev.ltms.fleet.msg.ReplyPushLoop -> dev.ltms.fleet.mcp.PrimaryRegistry",
|
||||
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.Rendezvous",
|
||||
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.Rendezvous$Resolution",
|
||||
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.TurnToken",
|
||||
"dev.ltms.fleet.inject.CompletionResolver$InFlight -> dev.ltms.fleet.msg.Rendezvous$Resolution",
|
||||
"dev.ltms.fleet.inject.Injector -> dev.ltms.fleet.msg.TurnToken",
|
||||
"dev.ltms.fleet.inject.Injector$Pending -> dev.ltms.fleet.msg.TurnToken",
|
||||
"dev.ltms.fleet.inject.TurnListener -> dev.ltms.fleet.msg.TurnToken",
|
||||
"dev.ltms.fleet.inject.TurnRegistrar -> dev.ltms.fleet.msg.TurnToken",
|
||||
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector",
|
||||
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector$Cancellation",
|
||||
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector$Delivery",
|
||||
"dev.ltms.fleet.metrics.FleetMetrics -> dev.ltms.fleet.msg.ReplyInbox",
|
||||
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.metrics.Metrics",
|
||||
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.metrics.Metrics",
|
||||
"dev.ltms.fleet.msg.ReplyPushLoop -> dev.ltms.fleet.metrics.Metrics",
|
||||
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.session.MemberSession",
|
||||
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.session.MemberSession$State",
|
||||
"dev.ltms.fleet.session.SessionManager -> dev.ltms.fleet.msg.TurnToken"
|
||||
);
|
||||
|
||||
private static final String ROOT_PACKAGE = "dev.ltms.fleet.";
|
||||
|
||||
@Test
|
||||
void packagesAreFreeOfCycles() {
|
||||
var classes = new ClassFileImporter()
|
||||
JavaClasses classes = new ClassFileImporter()
|
||||
.withImportOption(ImportOption.Predefined.DO_NOT_INCLUDE_TESTS)
|
||||
.importPackages("dev.ltms.fleet");
|
||||
|
||||
checkBaselineMatchesTodaysEdges(classes);
|
||||
|
||||
SliceRule rule = SlicesRuleDefinition.slices()
|
||||
.matching("dev.ltms.fleet.(*)..")
|
||||
.should().beFreeOfCycles();
|
||||
|
||||
// fleetd #131 step 1: move ConnectionIdentity so authz stops depending on the
|
||||
// MCP layer. Evidence: auth/CallerResolver.java:3 imports mcp.ConnectionIdentity;
|
||||
// mcp/FleetMcp.java:3-7 imports auth.AuditLog, Authz, CallerResolver, Principal,
|
||||
// Role.
|
||||
rule = ignoreCycle(rule, "auth", "mcp");
|
||||
|
||||
// fleetd #131 step 2: PrimaryRegistry is used by loops in msg; move it, or put
|
||||
// an interface between msg and mcp. Evidence: msg/ReplyPushLoop.java:5 and
|
||||
// msg/LeadHeartbeatLoop.java:5 import mcp.PrimaryRegistry; mcp/FleetMcp.java:15-18
|
||||
// imports msg.LeadChannel, LeadMessage, MessageService, Rendezvous.
|
||||
rule = ignoreCycle(rule, "mcp", "msg");
|
||||
|
||||
// fleetd #131 -- found while implementing this test, NOT one of the ticket's
|
||||
// original three; it names its own follow-up step before removal. Evidence:
|
||||
// inject/CompletionResolver.java:4-5, inject/Injector.java:6 and
|
||||
// inject/TurnListener.java:3 import msg.Rendezvous / msg.TurnToken;
|
||||
// msg/MessageService.java:6 imports inject.Injector.
|
||||
rule = ignoreCycle(rule, "inject", "msg");
|
||||
|
||||
// fleetd #131 -- same as above, its own follow-up. Evidence:
|
||||
// metrics/FleetMetrics.java:3 imports msg.ReplyInbox; msg/MessageService.java:7-8,
|
||||
// msg/LeadHeartbeatLoop.java:6-7 and msg/ReplyPushLoop.java:6-7 import
|
||||
// metrics.FleetMetrics / metrics.Metrics.
|
||||
rule = ignoreCycle(rule, "metrics", "msg");
|
||||
|
||||
// fleetd #131 -- same as above, its own follow-up. Evidence:
|
||||
// session/SessionManager.java:7 imports msg.TurnToken;
|
||||
// msg/LeadHeartbeatLoop.java:8 imports session.MemberSession.
|
||||
rule = ignoreCycle(rule, "msg", "session");
|
||||
|
||||
for (String edge : BASELINE_EDGES) {
|
||||
String[] originAndTarget = edge.split(" -> ");
|
||||
rule = rule.ignoreDependency(originAndTarget[0], originAndTarget[1]);
|
||||
}
|
||||
rule.check(classes);
|
||||
}
|
||||
|
||||
/**
|
||||
* Accepts today's known cycle between two top-level packages, and nothing else.
|
||||
* Ignoring both directions removes exactly this pair from cycle detection; every
|
||||
* other dependency -- including any new one added later, between these same two
|
||||
* packages or any other pair -- is still checked.
|
||||
* Fails with the exact offending edge when the live code and {@link #BASELINE_EDGES}
|
||||
* disagree: a dependency crossing a baselined package pair that is not in the baseline,
|
||||
* or a baseline entry whose dependency no longer exists.
|
||||
*/
|
||||
private static SliceRule ignoreCycle(SliceRule rule, String packageA, String packageB) {
|
||||
return rule
|
||||
.ignoreDependency(residesIn(packageA), residesIn(packageB))
|
||||
.ignoreDependency(residesIn(packageB), residesIn(packageA));
|
||||
private static void checkBaselineMatchesTodaysEdges(JavaClasses classes) {
|
||||
Set<String> baselinedPackagePairs = new TreeSet<>();
|
||||
for (String edge : BASELINE_EDGES) {
|
||||
String[] originAndTarget = edge.split(" -> ");
|
||||
baselinedPackagePairs.add(unorderedPair(
|
||||
topLevelPackageOf(originAndTarget[0]), topLevelPackageOf(originAndTarget[1])));
|
||||
}
|
||||
|
||||
Set<String> liveEdgesInBaselinedPairs = new TreeSet<>();
|
||||
for (JavaClass javaClass : classes) {
|
||||
for (Dependency dependency : javaClass.getDirectDependenciesFromSelf()) {
|
||||
JavaClass origin = dependency.getOriginClass();
|
||||
JavaClass target = dependency.getTargetClass();
|
||||
String originPackage = topLevelPackageOf(origin.getFullName());
|
||||
String targetPackage = topLevelPackageOf(target.getFullName());
|
||||
if (originPackage.isEmpty() || targetPackage.isEmpty() || originPackage.equals(targetPackage)) {
|
||||
continue;
|
||||
}
|
||||
if (baselinedPackagePairs.contains(unorderedPair(originPackage, targetPackage))) {
|
||||
liveEdgesInBaselinedPairs.add(origin.getFullName() + " -> " + target.getFullName());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
List<String> problems = new ArrayList<>();
|
||||
for (String liveEdge : liveEdgesInBaselinedPairs) {
|
||||
if (!BASELINE_EDGES.contains(liveEdge)) {
|
||||
String[] originAndTarget = liveEdge.split(" -> ");
|
||||
problems.add("new dependency not in the baseline: " + liveEdge
|
||||
+ " (packages " + topLevelPackageOf(originAndTarget[0])
|
||||
+ " -> " + topLevelPackageOf(originAndTarget[1]) + ")");
|
||||
}
|
||||
}
|
||||
for (String baselineEdge : BASELINE_EDGES) {
|
||||
if (!liveEdgesInBaselinedPairs.contains(baselineEdge)) {
|
||||
String[] originAndTarget = baselineEdge.split(" -> ");
|
||||
problems.add("stale baseline entry, no such dependency exists: " + baselineEdge
|
||||
+ " (packages " + topLevelPackageOf(originAndTarget[0])
|
||||
+ " -> " + topLevelPackageOf(originAndTarget[1]) + ")");
|
||||
}
|
||||
}
|
||||
|
||||
if (!problems.isEmpty()) {
|
||||
fail("PackageCyclesTest baseline is out of date:\n " + String.join("\n ", problems));
|
||||
}
|
||||
}
|
||||
|
||||
private static DescribedPredicate<JavaClass> residesIn(String topLevelPackage) {
|
||||
return Predicates.resideInAPackage("dev.ltms.fleet." + topLevelPackage + "..");
|
||||
private static String unorderedPair(String packageA, String packageB) {
|
||||
return packageA.compareTo(packageB) <= 0 ? packageA + "|" + packageB : packageB + "|" + packageA;
|
||||
}
|
||||
|
||||
private static String topLevelPackageOf(String fullyQualifiedClassName) {
|
||||
if (!fullyQualifiedClassName.startsWith(ROOT_PACKAGE)) {
|
||||
return "";
|
||||
}
|
||||
String rest = fullyQualifiedClassName.substring(ROOT_PACKAGE.length());
|
||||
int dot = rest.indexOf('.');
|
||||
return dot < 0 ? "" : rest.substring(0, dot);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -288,14 +288,16 @@ class AuthzTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* Every action beyond READ/METRICS/REPLY/ASK, asserted denied for an observer — including
|
||||
* {@code TASK_READ}, which is the entire point of this role: an unconfigured pane must not be
|
||||
* able to poll a ticket or read another session's status.
|
||||
* Every action beyond READ/METRICS/REPLY/ASK/SEND, asserted denied for an observer —
|
||||
* including {@code TASK_READ}, which is the entire point of this role: an unconfigured pane
|
||||
* must not be able to poll a ticket or read another session's status. {@code SEND} is excluded
|
||||
* here and given its own matrix below, since — unlike every action in this loop — its grant is
|
||||
* conditional on the target, not fixed.
|
||||
*/
|
||||
@Test
|
||||
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAndAsk() {
|
||||
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskAndSend() {
|
||||
for (Authz.Action a : Authz.Action.values()) {
|
||||
if (a == READ || a == METRICS || a == REPLY || a == ASK) {
|
||||
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == SEND) {
|
||||
continue;
|
||||
}
|
||||
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
|
||||
@@ -312,4 +314,53 @@ class AuthzTest {
|
||||
assertFalse(OBSERVER.isSpawnedMember());
|
||||
assertTrue(OBSERVER.isObserver());
|
||||
}
|
||||
|
||||
// ── the observer SEND matrix ────────────────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* {@code SEND} for an observer is the one grant that is conditional rather than fixed, exactly
|
||||
* like a collaborator's: flipping only the observer-target classifier's answer flips only this
|
||||
* outcome.
|
||||
*/
|
||||
@Test
|
||||
void anObserverMaySendOnlyWhenTheClassifierAcceptsTheTargetAsAnObserver() {
|
||||
assertTrue(Authz.permits(OBSERVER, SEND, "term_other_observer",
|
||||
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, target -> true),
|
||||
"the classifier accepting the target as an observer must grant SEND");
|
||||
assertFalse(Authz.permits(OBSERVER, SEND, "term_other_observer",
|
||||
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, target -> false),
|
||||
"the classifier refusing the target must deny SEND");
|
||||
assertFalse(Authz.permits(OBSERVER, SEND, "term_other_observer"),
|
||||
"the real production classifier recognises no terminal as an observer target yet, "
|
||||
+ "so SEND is refused by the three- and four-argument convenience forms");
|
||||
}
|
||||
|
||||
/**
|
||||
* Control for the test above: every other action's result for an observer does not move when
|
||||
* the observer-target classifier does. Only {@code SEND} is wired to it.
|
||||
*/
|
||||
@Test
|
||||
void theObserverTargetClassifierMovesOnlySendForAnObserver() {
|
||||
for (Authz.Action a : Authz.Action.values()) {
|
||||
if (a == SEND) {
|
||||
continue;
|
||||
}
|
||||
assertEquals(
|
||||
Authz.permits(OBSERVER, a, "term_observer"),
|
||||
Authz.permits(OBSERVER, a, "term_observer",
|
||||
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, target -> true),
|
||||
a + " must not depend on the observer-target classifier at all");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* An observer's {@code SEND} is gated on a different classifier than a collaborator's: the
|
||||
* collaborator classifier accepting every target must not itself grant an observer's SEND.
|
||||
*/
|
||||
@Test
|
||||
void anObserversSendDoesNotMoveOnTheCollaboratorClassifier() {
|
||||
assertFalse(Authz.permits(OBSERVER, SEND, "term_lead", target -> true),
|
||||
"an observer's SEND must consult the observer-target classifier, never the "
|
||||
+ "lead-or-collaborator one");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -912,4 +912,72 @@ class CallerResolverTest {
|
||||
assertFalse(r.knownLeadOrCollaborator().test("term_a"));
|
||||
}
|
||||
|
||||
// ── fleetd #743: sendableObserverTarget() reads the same maps and functions resolve() does ────
|
||||
|
||||
/**
|
||||
* A terminal this resolver recognises as none of the privileged roles is exactly the one
|
||||
* {@code resolve} would itself hand back {@link Role#OBSERVER} for.
|
||||
*/
|
||||
@Test
|
||||
void sendableObserverTargetIsTrueForATerminalKnownAsNoOtherRole() {
|
||||
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
|
||||
() -> Map.of("term_lead", "opus-5.0"), new MemberRegistry(null),
|
||||
t -> "term_worker".equals(t) ? MemberRole.DEV : null,
|
||||
() -> Map.of("term_collab", "ops"));
|
||||
|
||||
assertTrue(r.sendableObserverTarget().test("term_other"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendableObserverTargetIsFalseForALeadTerminal() {
|
||||
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
|
||||
() -> Map.of("term_lead", "opus-5.0"), new MemberRegistry(null), t -> null, Map::of);
|
||||
|
||||
assertFalse(r.sendableObserverTarget().test("term_lead"),
|
||||
"a lead's own terminal must never be a sendable observer target");
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendableObserverTargetIsFalseForACollaboratorTerminal() {
|
||||
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
|
||||
new MemberRegistry(null), t -> null, () -> Map.of("term_collab", "ops"));
|
||||
|
||||
assertFalse(r.sendableObserverTarget().test("term_collab"),
|
||||
"a collaborator's own terminal must never be a sendable observer target");
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendableObserverTargetIsFalseForALiveSpawnedMembersTerminal() {
|
||||
// Covers both a worker and an architect: spawnedMemberRole.apply(target) is non-null for
|
||||
// either, and resolve() never falls through to OBSERVER once it is.
|
||||
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
|
||||
new MemberRegistry(null),
|
||||
t -> switch (t) {
|
||||
case "term_worker" -> MemberRole.DEV;
|
||||
case "term_architect" -> MemberRole.ARCHITECT;
|
||||
default -> null;
|
||||
}, Map::of);
|
||||
|
||||
assertFalse(r.sendableObserverTarget().test("term_worker"));
|
||||
assertFalse(r.sendableObserverTarget().test("term_architect"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendableObserverTargetIsFalseForABoundArchitectSlotWithNoLiveMember() {
|
||||
// The edge case resolve() itself carries: a terminal bound to a configured architect slot
|
||||
// but with no live spawned-member session yet.
|
||||
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
|
||||
boundMembers("architect:lead-designer", MemberRole.ARCHITECT), t -> null, Map::of);
|
||||
|
||||
assertFalse(r.sendableObserverTarget().test("term_a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendableObserverTargetIsFalseForANullTarget() {
|
||||
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
|
||||
new MemberRegistry(null), t -> null, Map::of);
|
||||
|
||||
assertFalse(r.sendableObserverTarget().test(null));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -24,6 +24,41 @@ public final class FakeHerdr implements HerdrClient {
|
||||
/** The foreground PID of the one agent pane (term_a) in the canned {@code pane.process_info}. */
|
||||
public static final long WORKER_PID = 4242;
|
||||
|
||||
/**
|
||||
* A {@code detection} region of the current Claude Code TUI, whose input box is empty: a caret
|
||||
* line between two rules, above the footer.
|
||||
*/
|
||||
public static final String IDLE_PROMPT_CARET = """
|
||||
──────────────────────── lead: opus ─
|
||||
❯
|
||||
─────────────────────────────────────
|
||||
lead: opus · Opus 5 (1M context) · ~/LTMS/claude-bridge
|
||||
⏵⏵ auto mode on (shift+tab to cycle) · ← 1 agent""";
|
||||
|
||||
/** The same region with the operator's unsubmitted line still at the caret. */
|
||||
public static final String DRAFTED_PROMPT_CARET = """
|
||||
──────────────────────── lead: opus ─
|
||||
❯ yes, send it to lead: opus
|
||||
─────────────────────────────────────
|
||||
lead: opus · Opus 5 (1M context) · ~/LTMS/claude-bridge
|
||||
⏵⏵ auto mode on (shift+tab to cycle) · ← 1 agent""";
|
||||
|
||||
/** A {@code detection} region of an older TUI, which drew a bordered box, with that box empty. */
|
||||
public static final String IDLE_PROMPT_BOX = """
|
||||
⏺ done
|
||||
╭────────────────────────╮
|
||||
│ > │
|
||||
╰────────────────────────╯
|
||||
⏵⏵ auto mode on""";
|
||||
|
||||
/** The older TUI's bordered box, still holding the operator's unsubmitted line. */
|
||||
public static final String DRAFTED_PROMPT_BOX = """
|
||||
⏺ done
|
||||
╭────────────────────────╮
|
||||
│ > fix the issue when I │
|
||||
╰────────────────────────╯
|
||||
⏵⏵ auto mode on""";
|
||||
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
/**
|
||||
* Thread-safe on purpose. Background loops — {@link dev.ltms.fleet.msg.ReplyPushLoop} and the
|
||||
@@ -48,12 +83,15 @@ public final class FakeHerdr implements HerdrClient {
|
||||
private final Map<String, String> processInfoErrorCodeFor = new ConcurrentHashMap<>();
|
||||
private String tabCloseErrorCode = null;
|
||||
private final Map<String, String> tabCloseErrorCodeFor = new ConcurrentHashMap<>();
|
||||
private String workspaceListErrorCode = null;
|
||||
private String agentSendErrorCode = null;
|
||||
private boolean noPanes = false;
|
||||
private volatile String agentStatus = "idle"; // steady-state agent.get status
|
||||
private volatile String agentType = "claude"; // detected agent kind on agent.get; null = undetected
|
||||
private volatile String agentSessionId = null; // agent_session.value on agent.get; null = omitted
|
||||
private volatile String readText = "worker transcript tail"; // canned agent.read output
|
||||
/** Canned {@code detection}-source output, or {@code null} to serve {@link #readText} there too. */
|
||||
private volatile String detectionText = null;
|
||||
private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane
|
||||
private String pinnedStartTerminal;
|
||||
private String pinnedStartPane;
|
||||
@@ -142,6 +180,12 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make {@code workspace.list} fail with this herdr error code; every other method still succeeds. */
|
||||
public FakeHerdr workspaceListFailsWith(String code) {
|
||||
this.workspaceListErrorCode = code;
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code pane.list} report no panes at all — models a second herdr daemon (CB-185) that
|
||||
* simply does not host the pane a {@link PaneLocator} is searching for.
|
||||
@@ -200,12 +244,26 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/** The text {@code agent.read} returns (the CB-106 completion scrape). */
|
||||
/**
|
||||
* The text {@code agent.read} returns (the CB-106 completion scrape). It serves the
|
||||
* {@code detection} source as well unless {@link #detectionText} overrides that one.
|
||||
*/
|
||||
public FakeHerdr readText(String text) {
|
||||
this.readText = text;
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Override the text {@code agent.read} returns for the {@code detection} source only — the
|
||||
* prompt/footer tail herdr uses for status detection, a different region from the transcript the
|
||||
* other sources carry. Needed by a test whose subject reads the input box, since one
|
||||
* {@link #readText} cannot be both a worker's transcript and a lead's empty prompt.
|
||||
*/
|
||||
public FakeHerdr detectionText(String text) {
|
||||
this.detectionText = text;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make delivery ({@code agent.prompt} / {@code agent.send_keys}) fail with this error code. */
|
||||
public FakeHerdr agentSendFailsWith(String code) {
|
||||
this.agentSendErrorCode = code;
|
||||
@@ -302,11 +360,17 @@ public final class FakeHerdr implements HerdrClient {
|
||||
case "ping" -> mapper.readTree(
|
||||
("{\"type\":\"pong\",\"version\":\"%s\",\"protocol\":%d}")
|
||||
.formatted(pingVersion, pingProtocol));
|
||||
case "workspace.list" -> mapper.readTree(("""
|
||||
case "workspace.list" -> {
|
||||
if (workspaceListErrorCode != null) {
|
||||
throw new HerdrException("herdr error [" + workspaceListErrorCode + "]: workspace.list failed",
|
||||
workspaceListErrorCode, null);
|
||||
}
|
||||
yield mapper.readTree(("""
|
||||
{"type":"workspace_list","workspaces":[
|
||||
{"workspace_id":"w1","label":"dev-mgnl","focused":true,"pane_count":7,"agent_status":"unknown"},
|
||||
{"workspace_id":"w2","label":"ltms","focused":false,"pane_count":5,"agent_status":"done"}%s]}""")
|
||||
.formatted(extraWorkspaces.isEmpty() ? "" : "," + String.join(",", extraWorkspaces)));
|
||||
}
|
||||
case "agent.list" -> mapper.readTree(("""
|
||||
{"type":"agent_list","agents":[
|
||||
{"terminal_id":"term_a","agent":"claude","agent_status":"idle",
|
||||
@@ -347,8 +411,13 @@ public final class FakeHerdr implements HerdrClient {
|
||||
"agent_status":"%s","workspace_id":"w2","tab_id":"w2:t7","pane_id":"w2:p7"%s}}""")
|
||||
.formatted(agentField, agentStatus, sessionField));
|
||||
}
|
||||
case "agent.read" -> mapper.readTree(mapper.writeValueAsString(
|
||||
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", readText))));
|
||||
case "agent.read" -> {
|
||||
Object source = params instanceof Map<?, ?> m ? m.get("source") : null;
|
||||
String text = "detection".equals(source) && detectionText != null
|
||||
? detectionText : readText;
|
||||
yield mapper.readTree(mapper.writeValueAsString(
|
||||
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", text))));
|
||||
}
|
||||
case "agent.start" -> {
|
||||
if (onAgentStart != null) {
|
||||
onAgentStart.run();
|
||||
|
||||
@@ -0,0 +1,132 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/** Reading a Claude Code input box, so a paste-and-submit delivery never submits the operator's draft. */
|
||||
class PromptBoxTest {
|
||||
|
||||
private static final String EMPTY = FakeHerdr.IDLE_PROMPT_CARET;
|
||||
private static final String DRAFTED = FakeHerdr.DRAFTED_PROMPT_CARET;
|
||||
|
||||
// --- pure classification -------------------------------------------------
|
||||
|
||||
@Test
|
||||
void anEmptyCaretLineIsAnEmptyBox() {
|
||||
assertEquals(new PromptBox.Reading(PromptBox.State.EMPTY, 0), PromptBox.classify(EMPTY));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aCaretLineHoldingTextIsADraftAndCountsItsCharacters() {
|
||||
PromptBox.Reading reading = PromptBox.classify(DRAFTED);
|
||||
assertEquals(PromptBox.State.DRAFT, reading.state());
|
||||
assertEquals("yes,sendittolead:opus".length(), reading.characters(),
|
||||
"padding does not count — only what the operator typed");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBorderedBoxIsReadToo() {
|
||||
assertEquals(PromptBox.State.EMPTY, PromptBox.classify(FakeHerdr.IDLE_PROMPT_BOX).state(),
|
||||
"an older TUI draws a bordered box, and its panes must still be readable");
|
||||
assertEquals(PromptBox.State.DRAFT, PromptBox.classify(FakeHerdr.DRAFTED_PROMPT_BOX).state());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSingleTypedCharacterIsADraft() {
|
||||
assertEquals(PromptBox.State.DRAFT, PromptBox.classify("❯ f").state());
|
||||
assertEquals(PromptBox.State.DRAFT, PromptBox.classify("│ > f │").state());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aCursorBlockInAnOtherwiseEmptyBoxIsEmpty() {
|
||||
assertEquals(PromptBox.State.EMPTY, PromptBox.classify("❯ █").state(),
|
||||
"a terminal capture may leave the cursor cell in an empty box");
|
||||
assertEquals(PromptBox.State.EMPTY, PromptBox.classify("│ > █ │").state());
|
||||
}
|
||||
|
||||
@Test
|
||||
void theBoxLineIsReadToItsEndWhateverFollowsIt() {
|
||||
assertEquals(PromptBox.State.DRAFT, PromptBox.classify("❯ half a line\n ⏵⏵ auto mode on").state());
|
||||
assertEquals(PromptBox.State.EMPTY, PromptBox.classify("❯\n ⏵⏵ auto mode on").state());
|
||||
}
|
||||
|
||||
@Test
|
||||
void theLastBoxLineOnThePaneIsTheLiveOne() {
|
||||
assertEquals(PromptBox.State.DRAFT,
|
||||
PromptBox.classify("❯ an earlier prompt\n⏺ its answer\n❯ typing now").state(),
|
||||
"the detection region carries scrollback, so earlier prompts sit above the live box");
|
||||
assertEquals(PromptBox.State.EMPTY,
|
||||
PromptBox.classify("❯ an earlier prompt\n⏺ its answer\n❯").state());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aMarkerPartWayAlongALineIsNotABox() {
|
||||
assertEquals(PromptBox.State.UNREADABLE, PromptBox.classify("⏺ type ❯ to get a prompt").state(),
|
||||
"a caret the operator quoted is transcript text, not an input box");
|
||||
assertEquals(PromptBox.State.EMPTY, PromptBox.classify("⏺ type ❯ to get a prompt\n❯").state(),
|
||||
"and it must not shadow the real box further down");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aPaneWithNoBoxIsUnreadable() {
|
||||
assertEquals(PromptBox.State.UNREADABLE, PromptBox.classify("garbled ansi noise").state());
|
||||
assertEquals(PromptBox.State.UNREADABLE, PromptBox.classify("").state());
|
||||
assertEquals(PromptBox.State.UNREADABLE, PromptBox.classify(null).state());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aGeneratingTurnIsUnreadableEvenWithAnEmptyBox() {
|
||||
assertEquals(PromptBox.State.UNREADABLE,
|
||||
PromptBox.classify(EMPTY + "\n ✳ Thinking… (12s · esc to interrupt)").state());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aGeneratingMarkerInScrollbackAboveTheBoxDoesNotMakeThePaneUnreadable() {
|
||||
assertEquals(PromptBox.State.EMPTY,
|
||||
PromptBox.classify(" ✳ Thinking… (12s · esc to interrupt)\n⏺ done\n" + EMPTY).state(),
|
||||
"that marker survives in scrollback, and holding on it would hold every delivery forever");
|
||||
}
|
||||
|
||||
// --- the gate ------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void anEmptyBoxClearsTheGateAndReadsTheDetectionRegion() {
|
||||
FakeHerdr herdr = new FakeHerdr().detectionText(EMPTY);
|
||||
|
||||
assertTrue(new PromptBox(new AgentControl(herdr)).clearToSubmit("term_a"));
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
var params = (java.util.Map<String, Object>) herdr.lastCall("agent.read").params();
|
||||
assertEquals("detection", params.get("source"),
|
||||
"the input box is drawn in the detection region, not in transcript scrollback");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aDraftedBoxHoldsTheGate() {
|
||||
assertFalse(new PromptBox(new AgentControl(new FakeHerdr().detectionText(DRAFTED))).clearToSubmit("term_a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void anUnreadablePaneHoldsTheGate() {
|
||||
assertFalse(new PromptBox(new AgentControl(new FakeHerdr().detectionText("garbled"))).clearToSubmit("term_a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFailedReadHoldsTheGate() {
|
||||
assertFalse(new PromptBox(new AgentControl(new FakeHerdr().healthy(false))).clearToSubmit("term_a"),
|
||||
"a pane this cannot read must never be pasted into");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theGateClearsAgainOnceTheBoxEmpties() {
|
||||
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
|
||||
PromptBox box = new PromptBox(new AgentControl(herdr));
|
||||
|
||||
assertFalse(box.clearToSubmit("term_a"));
|
||||
herdr.detectionText(EMPTY);
|
||||
assertTrue(box.clearToSubmit("term_a"));
|
||||
}
|
||||
}
|
||||
@@ -89,6 +89,11 @@ class BackendOutageFlowTest {
|
||||
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", "term_primary").put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
// The lead-nudge paths read the input box before pasting into it.
|
||||
return MAPPER.createObjectNode().set("read",
|
||||
MAPPER.createObjectNode().put("text", FakeHerdr.IDLE_PROMPT_CARET));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
prompts.add(params);
|
||||
sendLatch.countDown();
|
||||
|
||||
@@ -59,6 +59,17 @@ class InjectorTest {
|
||||
assertEquals(List.of("hello"), sent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void deliveringToAMemberReadsNoPane() {
|
||||
// No human types into a spawned member's pane, so its delivery path must not pay for a
|
||||
// prompt-box read the way a lead's nudge paths do.
|
||||
injector.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertEquals(List.of("task"), sent());
|
||||
assertFalse(herdr.called("agent.read"), "a member delivery must not read its pane");
|
||||
}
|
||||
|
||||
@Test
|
||||
void holdsDeliveryUntilTheWorkerIsAvailable() {
|
||||
// CB-113: idle alone is not enough — hold until the worker's MCP is connected (ready).
|
||||
|
||||
@@ -16,6 +16,7 @@ import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -1127,8 +1128,13 @@ class LeadRolloverTest {
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
WorkspaceControl spaces = new WorkspaceControl(herdr);
|
||||
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
|
||||
// Only LEAD resolves to a configured name — every other terminal (e.g. the distinct
|
||||
// term_cap_N terminals evictionCountsInProgressEntriesTowardTheCap opens) falls back to
|
||||
// keying its own single-flight claim on the terminal itself, exactly like a terminal
|
||||
// the live roster does not recognise.
|
||||
return new LeadRollover(agents, spaces, launcher, () -> config, _ -> null,
|
||||
_ -> LEAD_NAME, () -> liveLeadTerminals, nowMillis, () -> { }, runner);
|
||||
t -> LEAD.equals(t) ? LEAD_NAME : null, () -> liveLeadTerminals, nowMillis, () -> { },
|
||||
runner);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -1317,7 +1323,7 @@ class LeadRolloverTest {
|
||||
+ "so an operator reading status() has something to act on: " + status.detail());
|
||||
}
|
||||
|
||||
// ---- fleetd #726 unit 3: confirm() single-flights per lead terminal -----------------------
|
||||
// ---- fleetd #726 unit 3: confirm() single-flights one roll at a time per lead ---------------
|
||||
|
||||
@Test
|
||||
@DisplayName("[SINGLE-FLIGHT 1] a second confirm() for the SAME lead terminal is refused with "
|
||||
@@ -1343,8 +1349,8 @@ class LeadRolloverTest {
|
||||
|
||||
assertFalse(secondDecision.accepted());
|
||||
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
|
||||
assertTrue(secondDecision.detail().contains(LEAD), "the refusal detail must name the lead "
|
||||
+ "terminal: " + secondDecision.detail());
|
||||
assertTrue(secondDecision.detail().contains(LEAD_NAME), "the refusal detail must name "
|
||||
+ "the single-flight key (the lead's name): " + secondDecision.detail());
|
||||
assertTrue(secondDecision.detail().contains(first.token()), "the refusal detail must name "
|
||||
+ "the token holding the claim: " + secondDecision.detail());
|
||||
assertEquals(1, runner.heldCount(), "the refused confirm() must never have reached "
|
||||
@@ -1502,4 +1508,160 @@ class LeadRolloverTest {
|
||||
+ "this lead terminal could never be rolled again: " + retryDecision.reason() + " / "
|
||||
+ retryDecision.detail());
|
||||
}
|
||||
|
||||
// ---- fleetd #737 unit 4: confirm() single-flights per LEAD NAME, not per lead terminal ------
|
||||
|
||||
@Test
|
||||
@DisplayName("[NAME-KEYED 1] a second confirm() for the SAME lead is refused with "
|
||||
+ "ROLL_ALREADY_RUNNING even when it is opened from a DIFFERENT terminal, while the "
|
||||
+ "first roll's continuation is still in flight")
|
||||
void secondConfirmForTheSameLeadIsRefusedEvenFromADifferentTerminal() throws IOException {
|
||||
FakeHerdr herdr = herdrReadyForAFullRoll(); // the held roll WOULD complete once run
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
HoldingRunner runner = new HoldingRunner();
|
||||
Map<String, String> liveLeadTerminals = Map.of(NEW_TERMINAL, LEAD_NAME);
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
WorkspaceControl spaces = new WorkspaceControl(herdr);
|
||||
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
|
||||
// Both LEAD and OTHER_LEAD resolve to the SAME configured lead name — modelling a roll
|
||||
// that has already replaced the lead's pane: the fresh pane (here, OTHER_LEAD) is still
|
||||
// the SAME lead, just a different terminal id.
|
||||
Function<String, String> leadNameForTerminal = _ -> LEAD_NAME;
|
||||
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
|
||||
_ -> null, leadNameForTerminal, () -> liveLeadTerminals, fixedClock(clock), () -> { },
|
||||
runner);
|
||||
|
||||
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision firstDecision = rollover.confirm(LEAD, first.token(), true);
|
||||
assertTrue(firstDecision.accepted(), "expected approval; got: " + firstDecision.reason()
|
||||
+ " / " + firstDecision.detail());
|
||||
assertEquals(1, runner.heldCount(), "sanity: the first roll is held, not run yet");
|
||||
|
||||
LeadRollover.PendingRollover second = rollover.open(OTHER_LEAD,
|
||||
"a second request for the SAME lead, opened from a DIFFERENT terminal");
|
||||
LeadRollover.RollDecision secondDecision = rollover.confirm(OTHER_LEAD, second.token(), true);
|
||||
|
||||
assertFalse(secondDecision.accepted(), "a different terminal resolving to the SAME lead "
|
||||
+ "name must still be refused — the single-flight claim is keyed on the lead's "
|
||||
+ "name, not its terminal");
|
||||
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
|
||||
assertTrue(secondDecision.detail().contains(LEAD_NAME), "the refusal detail must name "
|
||||
+ "the single-flight key (the lead's name): " + secondDecision.detail());
|
||||
assertTrue(secondDecision.detail().contains(first.token()), "the refusal detail must "
|
||||
+ "name the token holding the claim: " + secondDecision.detail());
|
||||
assertEquals(1, runner.heldCount(), "the refused confirm() must never have reached "
|
||||
+ "continuationRunner — it ran exactly once, for the first roll only");
|
||||
|
||||
runner.runNext(); // let the first (and only) held roll finish
|
||||
|
||||
// The claim must have been released once the first roll's continuation finished — a
|
||||
// fresh request for the SAME lead, even opened from yet another terminal, is now
|
||||
// approved.
|
||||
LeadRollover.PendingRollover third = rollover.open(OTHER_LEAD, "retry after the first roll finished");
|
||||
LeadRollover.RollDecision thirdDecision = rollover.confirm(OTHER_LEAD, third.token(), true);
|
||||
assertTrue(thirdDecision.accepted(), "the claim must have been released once the first "
|
||||
+ "roll finished: " + thirdDecision.reason() + " / " + thirdDecision.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[NAME-KEYED 2] the claim is released under the key it was taken under, even "
|
||||
+ "when the old terminal has already dropped out of the live-lead roster by release "
|
||||
+ "time — success (non-throwing) path")
|
||||
void claimReleasedUnderCarriedKeyEvenWhenOldTerminalIsNoLongerInTheRosterSuccessPath() throws IOException {
|
||||
// herdrReadyForAFullRoll() lets the RETRY below run its continuation to a genuine ROLLED
|
||||
// completion (pane gone on the first check, a pinned relaunch) — needed because, unlike
|
||||
// the other NAME-KEYED tests, this one's retry resolves a REAL name ("opus") and must not
|
||||
// spin forever against this test's fixed, non-advancing clock (see this class's javadoc).
|
||||
FakeHerdr herdr = herdrReadyForAFullRoll();
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
WorkspaceControl spaces = new WorkspaceControl(herdr);
|
||||
// A config that CAN relaunch 'opus' — unlike emptyFleetConfig(), its fleet() is non-null,
|
||||
// so resolveLaunchable(null) below returns null cleanly instead of throwing a NullPointer
|
||||
// out of cfg.fleet() itself; this test needs the NON-throwing relaunch-refused exit, not
|
||||
// an incidental NPE.
|
||||
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
|
||||
Map<String, String> roster = new HashMap<>();
|
||||
roster.put(LEAD, LEAD_NAME);
|
||||
Map<String, String> liveLeadTerminals = Map.of(NEW_TERMINAL, LEAD_NAME);
|
||||
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
|
||||
_ -> null, roster::get, () -> liveLeadTerminals, fixedClock(clock), () -> { }, Runnable::run);
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
|
||||
// The live roster drops the OLD terminal before this roll's continuation runs — the
|
||||
// hazard fleetd #737 names: leadNameForTerminal reads the LIVE roster, and (with the
|
||||
// synchronous runner this test injects) the continuation runs INSIDE this confirm()
|
||||
// call, strictly after open() already resolved and carried the key.
|
||||
roster.remove(LEAD);
|
||||
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.RELAUNCH_FAILED, status.state(), "sanity: leadName "
|
||||
+ "resolved to null (the roster had already dropped the old terminal), so "
|
||||
+ "relaunch(null) fails cleanly, through the non-throwing exit this test means "
|
||||
+ "to cover: " + status.detail());
|
||||
|
||||
// A fresh terminal for the SAME lead name, confirmed live again — proves the claim this
|
||||
// roll took under the lead's name was actually released, not left stuck under whatever
|
||||
// a release that RECOMPUTES the key via leadNameForTerminal (now null for the old
|
||||
// terminal) would have tried to remove instead.
|
||||
roster.put(OTHER_LEAD, LEAD_NAME);
|
||||
LeadRollover.PendingRollover retry = rollover.open(OTHER_LEAD, "a fresh pane for the same lead");
|
||||
LeadRollover.RollDecision retryDecision = rollover.confirm(OTHER_LEAD, retry.token(), true);
|
||||
assertTrue(retryDecision.accepted(), "the claim taken under the lead's name must have "
|
||||
+ "been released even though the OLD terminal no longer resolved to any name at "
|
||||
+ "release time: " + retryDecision.reason() + " / " + retryDecision.detail());
|
||||
LeadRollover.RollStatus retryStatus = rollover.status(retry.token());
|
||||
assertEquals(LeadRollover.RollState.ROLLED, retryStatus.state(), "sanity: this retry's "
|
||||
+ "own claim (also 'opus') must not have been blocked by a leftover claim from "
|
||||
+ "the first roll: " + retryStatus.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[NAME-KEYED 3] the claim is released under the key it was taken under on the "
|
||||
+ "thrown-exception path too, even when the old terminal has already dropped out of "
|
||||
+ "the live-lead roster")
|
||||
void claimReleasedUnderCarriedKeyEvenWhenOldTerminalIsNoLongerInTheRosterThrowingPath() throws IOException {
|
||||
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
|
||||
fake.paneCloseFailsWith("permission_denied"); // a real failure, not an already-gone code
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
AgentControl agents = new AgentControl(fake);
|
||||
WorkspaceControl spaces = new WorkspaceControl(fake);
|
||||
LeadLauncher launcher = fakeLauncher(fake, emptyFleetConfig());
|
||||
Map<String, String> roster = new HashMap<>();
|
||||
roster.put(LEAD, LEAD_NAME);
|
||||
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
|
||||
_ -> null, roster::get, Map::of, fixedClock(clock), () -> { }, Runnable::run);
|
||||
|
||||
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
|
||||
|
||||
// Same hazard as [NAME-KEYED 2] above, now exercised on the thrown-exception exit.
|
||||
roster.remove(LEAD);
|
||||
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, first.token(), true);
|
||||
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
|
||||
+ "inside the deferred continuation, which this test's synchronous runner has "
|
||||
+ "already run to completion by the time confirm() returns");
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(first.token());
|
||||
assertEquals(LeadRollover.RollState.FAILED, status.state(), "sanity: the continuation "
|
||||
+ "threw and left a terminal FAILED outcome: " + status.detail());
|
||||
|
||||
// A fresh terminal for the SAME lead name, confirmed live again — proves the claim this
|
||||
// roll took under the lead's name was actually released, not left stuck under whatever
|
||||
// a release that RECOMPUTES the key via leadNameForTerminal (now null for the old
|
||||
// terminal) would have tried to remove instead.
|
||||
roster.put(OTHER_LEAD, LEAD_NAME);
|
||||
LeadRollover.PendingRollover retry = rollover.open(OTHER_LEAD, "retry after the throw");
|
||||
LeadRollover.RollDecision retryDecision = rollover.confirm(OTHER_LEAD, retry.token(), true);
|
||||
assertTrue(retryDecision.accepted(), "the claim must have been released even though the "
|
||||
+ "continuation threw AND the old terminal no longer resolved to any name at "
|
||||
+ "release time: " + retryDecision.reason() + " / " + retryDecision.detail());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@ import dev.ltms.fleet.metrics.Metrics;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.FakeWorktrees;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
@@ -269,6 +270,60 @@ class FleetMcpAuthzTest {
|
||||
"a spawned member's own terminal must stay unreachable, even once the classifier is real");
|
||||
}
|
||||
|
||||
// --- fleetd #743: the observer SEND matrix, over MCP's denyFor -------------------------------
|
||||
|
||||
private static final Principal OBSERVER = Principal.observer("term_observer", 700);
|
||||
|
||||
/**
|
||||
* Wires one real {@link CallerResolver} that recognises a lead, a collaborator, and a live
|
||||
* spawned worker, leaving "term_other_observer" classified as none of them — so the same
|
||||
* wiring both denies an observer's {@code SEND} to every privileged role and grants it to
|
||||
* another unclassified pane, proving the refusals are the rule and not a missing fixture.
|
||||
*/
|
||||
@Test
|
||||
void anObserverMaySendOnlyToAnotherObserverNeverToALeadWorkerOrArchitect() {
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
() -> Map.of("term_lead_known", "lead-x"), new MemberRegistry(null),
|
||||
t -> "term_a".equals(t) ? MemberRole.DEV : null,
|
||||
() -> Map.of("term_collab_known", "ops2"));
|
||||
FleetMcp m = mcp(true, callers);
|
||||
|
||||
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_lead_known"),
|
||||
"an observer must never reach a lead's terminal");
|
||||
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_collab_known"),
|
||||
"an observer must never reach a collaborator's terminal");
|
||||
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_a"),
|
||||
"an observer must never reach a live spawned member's terminal");
|
||||
|
||||
// CONTROL: the same wiring, the same denyFor call, a target recognised as none of the
|
||||
// three privileged roles above -- this is what proves the three refusals above are the
|
||||
// rule working, not a classifier that refuses every target regardless of what it is.
|
||||
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_other_observer"),
|
||||
"an observer must reach another pane that resolves as an observer itself");
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #anObserverMaySendOnlyToAnotherObserverNeverToALeadWorkerOrArchitect}, for a
|
||||
* terminal bound to a configured architect slot but hosting no live spawned-member session --
|
||||
* the case {@link CallerResolver#resolve} itself treats separately from a live worker/architect.
|
||||
*/
|
||||
@Test
|
||||
void anObserverMayNotSendToABoundArchitectSlotEither() {
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
|
||||
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(Map.of(),
|
||||
Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null));
|
||||
assertTrue(members.bind("architect:lead-designer", "term_bound_architect"));
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, Map::of,
|
||||
members, t -> null, Map::of);
|
||||
FleetMcp m = mcp(true, callers);
|
||||
|
||||
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_bound_architect"),
|
||||
"a terminal bound to a configured architect slot must stay unreachable to an observer");
|
||||
// CONTROL: the same wiring, a target the bind above never touched.
|
||||
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_other_observer"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void theLegacyConstructorLeavesTheGateOpen() {
|
||||
// The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced.
|
||||
@@ -453,6 +508,28 @@ class FleetMcpAuthzTest {
|
||||
assertFalse(FleetMcp.membersVisibleTo(ANON), "authenticated as nothing must not see it either");
|
||||
}
|
||||
|
||||
// --- who may see fleet_list's panes array ----------------------------------------------------
|
||||
|
||||
/**
|
||||
* {@link FleetMcp#panesVisibleTo} is the whole policy decision for {@code fleet_list}'s
|
||||
* {@code panes} array: visible to every role that may {@code SEND} to some other pane -- the
|
||||
* primary, an architect, a collaborator, and an observer (to another observer pane only, with
|
||||
* its row filtered and reduced -- see {@code listFleet}) -- never a worker, never an
|
||||
* anonymous caller.
|
||||
*/
|
||||
@Test
|
||||
void onlyPrimaryArchitectCollaboratorAndObserverMaySeeThePanesArray() {
|
||||
assertTrue(FleetMcp.panesVisibleTo(PRIMARY), "the primary must see the panes array");
|
||||
assertTrue(FleetMcp.panesVisibleTo(ARCH_DESIGN), "an architect must see the panes array");
|
||||
assertTrue(FleetMcp.panesVisibleTo(COLLABORATOR), "a collaborator must see its own peer roster");
|
||||
assertFalse(FleetMcp.panesVisibleTo(WORKER_A),
|
||||
"a worker holds READ but can never SEND, so it must not see the panes array");
|
||||
assertTrue(FleetMcp.panesVisibleTo(Principal.observer("term_obs", 700)),
|
||||
"an observer holds SEND to another observer pane, so it must see the (filtered, "
|
||||
+ "reduced) panes array");
|
||||
assertFalse(FleetMcp.panesVisibleTo(ANON), "authenticated as nothing must not see it either");
|
||||
}
|
||||
|
||||
/**
|
||||
* Same reasoning as {@link #theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}: the
|
||||
* predicate above can be perfectly correct while the one production call site never asks it.
|
||||
|
||||
@@ -0,0 +1,142 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.session.FakeWorktrees;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import io.modelcontextprotocol.client.McpClient;
|
||||
import io.modelcontextprotocol.client.McpSyncClient;
|
||||
import io.modelcontextprotocol.client.transport.HttpClientStreamableHttpTransport;
|
||||
import io.modelcontextprotocol.spec.McpClientTransport;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import org.eclipse.jetty.server.Server;
|
||||
import org.eclipse.jetty.server.ServerConnector;
|
||||
import org.eclipse.jetty.servlet.ServletContextHandler;
|
||||
import org.eclipse.jetty.servlet.ServletHolder;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #743, driven end to end: a real MCP client over a real HTTP transport, resolved by the
|
||||
* real {@link CallerResolver} to {@link dev.ltms.fleet.auth.Role#OBSERVER}, sending to another
|
||||
* unclassified pane. {@link FleetMcpAuthzTest} already proves {@code denyFor} grants this case and
|
||||
* that the handler calls {@code attributeIfObserver}; this test is the one path that also proves
|
||||
* the grant is not dead at a second gate (CB-548's two-gate trap) by driving the real
|
||||
* {@link Injector} to the point of its real herdr call, and reads the exact text the receiving
|
||||
* pane would see.
|
||||
*/
|
||||
class FleetMcpObserverSendDeliveryTest {
|
||||
|
||||
private static final String TARGET = "term_other_observer";
|
||||
|
||||
private final FakeHerdr herdr = new FakeHerdr();
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private final Injector injector = new Injector(agents);
|
||||
private final Rendezvous rendezvous = new Rendezvous();
|
||||
private final MessageService messages = new MessageService(agents, injector, rendezvous,
|
||||
new InMemoryReplyInbox());
|
||||
private FleetMcp mcp;
|
||||
private Server server;
|
||||
|
||||
@AfterEach
|
||||
void tearDown() throws Exception {
|
||||
if (server != null) {
|
||||
server.stop();
|
||||
}
|
||||
if (mcp != null) {
|
||||
mcp.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anObserversSendIsAttributedAndReachesTheRealInjector() throws Exception {
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(agents, new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
_ -> "tok");
|
||||
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
|
||||
// The fake's pane list carries a second pane, "term_shell", whose shell pid is 9001 and
|
||||
// which hosts no agent -- a herdr-owned pane recognised as no lead, architect, collaborator
|
||||
// or live spawned member, so the real resolver lands it on the observer floor.
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 9001L);
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, new MemberRegistry(null));
|
||||
|
||||
mcp = new FleetMcp(messages, workers, sessions, identity, sessions.asPresence(),
|
||||
new PrimaryRegistry(null), callers, FleetMcp.AuthorizationMode.ENFORCED,
|
||||
null, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), null, FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), List.of(), null);
|
||||
|
||||
ServletContextHandler handler = new ServletContextHandler();
|
||||
handler.setContextPath("/");
|
||||
handler.addServlet(new ServletHolder(mcp.servlet()), "/mcp");
|
||||
server = new Server(0);
|
||||
server.setHandler(handler);
|
||||
server.start();
|
||||
String baseUrl = "http://127.0.0.1:"
|
||||
+ ((ServerConnector) server.getConnectors()[0]).getLocalPort();
|
||||
|
||||
McpSchema.CallToolResult result = sendFleetSend(baseUrl, TARGET, "hi there");
|
||||
assertFalse(result.isError(), "an observer sending to another observer must be accepted: "
|
||||
+ textOf(result));
|
||||
|
||||
long waiterDeadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(TARGET) && System.currentTimeMillis() < waiterDeadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(TARGET), "the async send must have opened its rendezvous waiter");
|
||||
|
||||
injector.onStatus(TARGET, AgentStatus.IDLE); // drives the real delivery attempt to herdr
|
||||
|
||||
long deliveryDeadline = System.currentTimeMillis() + 3000;
|
||||
while (!herdr.called("agent.prompt") && System.currentTimeMillis() < deliveryDeadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(herdr.called("agent.prompt"), "the delivery attempt must have reached herdr");
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Object> params = (Map<String, Object>) herdr.lastCall("agent.prompt").params();
|
||||
assertEquals("[fleet_send from observer term_shell]\nhi there", params.get("text"),
|
||||
"the receiving pane must see the sender's own daemon-resolved terminal, never a raw "
|
||||
+ "echo of the content and never a client-supplied name");
|
||||
}
|
||||
|
||||
private static McpSchema.CallToolResult sendFleetSend(String baseUrl, String target, String content) {
|
||||
McpClientTransport transport = HttpClientStreamableHttpTransport.builder(baseUrl)
|
||||
.endpoint("/mcp")
|
||||
.build();
|
||||
try (McpSyncClient client = McpClient.sync(transport).build()) {
|
||||
client.initialize();
|
||||
return client.callTool(McpSchema.CallToolRequest.builder("fleet_send")
|
||||
.arguments(Map.of("sessionId", target, "content", content, "wait", false))
|
||||
.build());
|
||||
}
|
||||
}
|
||||
|
||||
private static String textOf(McpSchema.CallToolResult r) {
|
||||
return ((McpSchema.TextContent) r.content().getFirst()).text();
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.Fleetd;
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.auth.Principal;
|
||||
@@ -1110,6 +1111,215 @@ class FleetMcpTest {
|
||||
assertTrue(asArchitect.contains("\"sessionId\":\"term_collab\""), asArchitect);
|
||||
}
|
||||
|
||||
// --- fleet_list's panes array -----------------------------------------------------------------
|
||||
|
||||
/** Calls the canonical {@code listFleet} overload directly, so a test can set the pane-discovery
|
||||
* payload and its visibility independently of a real {@code Principal} / MCP exchange. */
|
||||
private static McpSchema.CallToolResult listFleetWithPanes(FakeHerdr h, FleetMcp.PaneSource panes,
|
||||
boolean panesVisible) {
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
return FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
|
||||
Map.of(), "", Map.of(), false,
|
||||
FleetMcp.CoordinationSource.none(), false, true, true, panes, panesVisible);
|
||||
}
|
||||
|
||||
/**
|
||||
* A pane whose tab herdr reports with a label gets that label and the exact terminal id
|
||||
* {@code fleet_send} takes as a target, carried as {@code sessionId}.
|
||||
*/
|
||||
@Test
|
||||
void listReportsAPaneRowWithItsTabLabelAndSendableSessionId() {
|
||||
FakeHerdr h = new FakeHerdr().withTab("w2", "w2:t7", "trinotes");
|
||||
MemberPresence presence = new MemberPresence();
|
||||
presence.markPresent("term_a");
|
||||
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
|
||||
() -> new PaneLocator(h).tabLabelsByTabId(),
|
||||
Fleetd.deliverableTo(presence, Map::of, Map::of), _ -> false, _ -> false);
|
||||
|
||||
String out = textOf(listFleetWithPanes(h, panes, true));
|
||||
|
||||
assertTrue(out.contains("\"panes\":["), out);
|
||||
assertTrue(out.contains("\"sessionId\":\"term_a\""), out);
|
||||
assertTrue(out.contains("\"label\":\"trinotes\""), out);
|
||||
assertTrue(out.contains("\"deliverable\":true"), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* A pane whose tab carries no label known to herdr still gets a row -- a missing label must
|
||||
* never throw, and must never drop the pane from the array, only report a {@code null} label.
|
||||
* Pairs with a {@code deliverable} false reading when the target is neither present, a lead,
|
||||
* nor a collaborator.
|
||||
*/
|
||||
@Test
|
||||
void listReportsAPaneRowWithANullLabelWhenHerdrHasNoneAndNotDeliverable() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
|
||||
() -> new PaneLocator(h).tabLabelsByTabId(),
|
||||
Fleetd.deliverableTo(new MemberPresence(), Map::of, Map::of), _ -> false, _ -> false);
|
||||
|
||||
String out = textOf(listFleetWithPanes(h, panes, true));
|
||||
|
||||
assertTrue(out.contains("\"panes\":["), out);
|
||||
assertTrue(out.contains("\"sessionId\":\"term_a\""), out);
|
||||
assertTrue(out.contains("\"label\":null"), out);
|
||||
assertTrue(out.contains("\"deliverable\":false"), out);
|
||||
}
|
||||
|
||||
/** A caller this role may not show the array to gets no {@code panes} key at all. */
|
||||
@Test
|
||||
void listOmitsThePanesArrayWhenTheCallerMayNotSeeIt() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
|
||||
String out = textOf(listFleetWithPanes(h, FleetMcp.PaneSource.none(), false));
|
||||
|
||||
assertFalse(out.contains("\"panes\""), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* The tab-label scan behind {@code panes} shares no failure path with the rest of
|
||||
* {@code listFleet} -- a {@code workspace.list}/{@code tab.list} failure costs only the
|
||||
* labels in the {@code panes} row (each renders {@code null}), never the {@code leads}/
|
||||
* {@code members} arrays, which never needed that scan at all.
|
||||
*/
|
||||
@Test
|
||||
void listStillReportsEveryOtherArrayWhenTheLabelScanFails() {
|
||||
FakeHerdr h = new FakeHerdr().workspaceListFailsWith("unavailable");
|
||||
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
|
||||
() -> new PaneLocator(h).tabLabelsByTabId(),
|
||||
Fleetd.deliverableTo(new MemberPresence(), Map::of, Map::of), _ -> false, _ -> false);
|
||||
|
||||
McpSchema.CallToolResult res = listFleetWithPanes(h, panes, true);
|
||||
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), textOf(res));
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"panes\":["), out);
|
||||
assertTrue(out.contains("\"sessionId\":\"term_a\""), out);
|
||||
assertTrue(out.contains("\"label\":null"), out);
|
||||
assertTrue(out.contains("\"leads\":[]"), "a label-scan failure must not cost the leads array: " + out);
|
||||
assertTrue(out.contains("\"members\":[]"), "a label-scan failure must not cost the members array: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #756: a pane bound to a configured architect slot with no live member session must
|
||||
* report {@code role: "architect"}, read from {@link CallerResolver#boundToArchitectSlot} —
|
||||
* the same classifier {@link CallerResolver#sendableObserverTarget} refuses as a {@code SEND}
|
||||
* target — rather than falling through to {@code "observer"}.
|
||||
*/
|
||||
@Test
|
||||
void listReportsArchitectForASlotBoundPaneWithNoLiveMember() {
|
||||
FakeHerdr h = new FakeHerdr()
|
||||
.withAgent("claude-arch", "term_unoccupied_architect", "w2:pArch", "w2:tArch");
|
||||
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
|
||||
Map::of, _ -> false, "term_unoccupied_architect"::equals, _ -> false);
|
||||
|
||||
String out = textOf(listFleetWithPanes(h, panes, true));
|
||||
|
||||
assertTrue(out.contains("\"sessionId\":\"term_unoccupied_architect\""), out);
|
||||
assertTrue(out.contains("\"role\":\"architect\""),
|
||||
"a slot-bound pane with no live session must read \"architect\", not the generic "
|
||||
+ "\"observer\" fallback: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* Calls the canonical {@code listFleet} overload directly with an explicit {@code leads}/
|
||||
* {@code collaborators} payload and {@code callerIsObserver}, mirroring exactly what the real
|
||||
* {@code fleet_list} handler computes for an observer caller: {@code panesVisible} true,
|
||||
* {@code leadsVisible}/{@code membersVisible}/{@code collaboratorsVisible} false.
|
||||
*/
|
||||
private static McpSchema.CallToolResult listFleetAsObserver(FakeHerdr h, Map<String, String> leads,
|
||||
Map<String, String> collaborators, FleetMcp.PaneSource panes) {
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
return FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
|
||||
leads, "", collaborators, false,
|
||||
FleetMcp.CoordinationSource.none(), false, false, false, panes, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #758: an observer's {@code fleet_list} now carries a {@code panes} key, filtered to
|
||||
* {@link CallerResolver#sendableObserverTarget} (so a lead's pane, a spawned member's pane, a
|
||||
* collaborator's pane, and an unoccupied architect-slot pane are all absent) and every
|
||||
* surviving row reduced to exactly {@code sessionId}, {@code label}, {@code status},
|
||||
* {@code role}, {@code deliverable} — never {@code paneId}, {@code workspaceId}, {@code tabId},
|
||||
* {@code agentType}, or {@code cwd}.
|
||||
*/
|
||||
@Test
|
||||
void listFiltersAndReducesThePanesArrayForAnObserver() {
|
||||
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(Map.of(),
|
||||
Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null));
|
||||
assertTrue(members.bind("architect:lead-designer", "term_architect_pane"));
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(null, false, null,
|
||||
() -> Map.of("term_lead_pane", "fleet01-lead"), members,
|
||||
t -> "term_member_pane".equals(t) ? MemberRole.DEV : null,
|
||||
() -> Map.of("term_collab_pane", "ops"));
|
||||
|
||||
FakeHerdr h = new FakeHerdr()
|
||||
.withAgent("claude-sendable", "term_sendable", "w2:pS", "w2:tS")
|
||||
.withAgent("claude-lead", "term_lead_pane", "w2:pL", "w2:tL")
|
||||
.withAgent("claude-member", "term_member_pane", "w2:pM", "w2:tM")
|
||||
.withAgent("claude-collab", "term_collab_pane", "w2:pC", "w2:tC")
|
||||
.withAgent("claude-arch", "term_architect_pane", "w2:pA", "w2:tA")
|
||||
.withTab("w2", "w2:tS", "trinotes");
|
||||
|
||||
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
|
||||
() -> new PaneLocator(h).tabLabelsByTabId(), _ -> true,
|
||||
callers::boundToArchitectSlot, callers.sendableObserverTarget());
|
||||
|
||||
String out = textOf(listFleetAsObserver(h, callers.leads(),
|
||||
Map.of("term_collab_pane", "ops"), panes));
|
||||
|
||||
assertTrue(out.contains("\"panes\":["), out);
|
||||
assertTrue(out.contains("\"sessionId\":\"term_sendable\""),
|
||||
"an ordinary unclassified pane must still be sendable and visible: " + out);
|
||||
assertFalse(out.contains("term_lead_pane"), "a lead's pane must not be enumerated: " + out);
|
||||
assertFalse(out.contains("term_member_pane"), "a spawned member's pane must not be enumerated: " + out);
|
||||
assertFalse(out.contains("term_collab_pane"), "a collaborator's pane must not be enumerated: " + out);
|
||||
assertFalse(out.contains("term_architect_pane"),
|
||||
"an unoccupied architect-slot pane must not be enumerated: " + out);
|
||||
assertFalse(out.contains("\"paneId\""), "an observer's row must never carry paneId: " + out);
|
||||
assertFalse(out.contains("\"workspaceId\""), "an observer's row must never carry workspaceId: " + out);
|
||||
assertFalse(out.contains("\"tabId\""), "an observer's row must never carry tabId: " + out);
|
||||
assertFalse(out.contains("\"agentType\""), "an observer's row must never carry agentType: " + out);
|
||||
assertFalse(out.contains("\"cwd\""), "an observer's row must never carry cwd: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* Control for the test above: a primary's {@code panes} row is unchanged by fleetd #758 —
|
||||
* {@code callerIsObserver} false keeps every field, including {@code paneId} and a spawned
|
||||
* member's {@code cwd}.
|
||||
*/
|
||||
@Test
|
||||
void listKeepsTheFullPaneRowForAPrimaryIncludingCwdAndPaneId() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
MemberSession spawned = sessions.acquire("ltms-local", "/worktree/member-1", null, null);
|
||||
h.withAgent("claude-member", spawned.terminalId(), "w9:pMember", "w9:tMember");
|
||||
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(Map::of, _ -> true, _ -> false, _ -> false);
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
|
||||
Map.of(), "", Map.of(), false,
|
||||
FleetMcp.CoordinationSource.none(), true, true, true, panes, true, false);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"sessionId\":\"" + spawned.terminalId() + "\""), out);
|
||||
assertTrue(out.contains("\"paneId\":\"w9:pMember\""),
|
||||
"a primary must still see the fleet_stop handle: " + out);
|
||||
assertTrue(out.contains("\"cwd\":\"/worktree/member-1\""),
|
||||
"a primary must still see a spawned member's worktree path: " + out);
|
||||
assertTrue(out.contains("\"role\":\"dev\""), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #421: {@code mailbox.pending} counts only broker-ready messages, so a blocked lead's
|
||||
* normal, healthy state is {@code "pending": 0} next to a non-empty {@code held[]} — which
|
||||
|
||||
@@ -40,7 +40,7 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void deliversAHeldMessageToTheLeadPaneAndAcksIt() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
|
||||
|
||||
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
|
||||
|
||||
@@ -57,7 +57,7 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void redeliveryOfAMessageAlreadyWrittenToThePaneIsAckedWithoutAnotherPaneWrite() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
|
||||
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
|
||||
|
||||
loop.tick();
|
||||
@@ -71,7 +71,7 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void aRedeliveryIsAckedEvenWhileTheLeadIsMidTurn() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
|
||||
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
|
||||
|
||||
loop.tick();
|
||||
@@ -91,7 +91,7 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void leavesTheMessageUnackedWhenTheLeadIsMidTurn() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
|
||||
var herdr = new FakeHerdr().agentStatus("working");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("working");
|
||||
|
||||
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
|
||||
|
||||
@@ -103,7 +103,7 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void leavesTheMessageUnackedWhenNoLeadPaneIsKnown() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
|
||||
|
||||
loop(channel, herdr, Map.of()).tick();
|
||||
|
||||
@@ -115,7 +115,8 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void leavesTheMessageUnackedWhenHerdrRefusesTheInjection() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("agent_not_found");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET)
|
||||
.agentStatus("idle").agentSendFailsWith("agent_not_found");
|
||||
|
||||
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
|
||||
|
||||
@@ -126,7 +127,7 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void resolvesTheLeadByNameWhenSeveralAreKnown() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
|
||||
// Two leads on this daemon; only one carries the coord-id the mailbox is owned as.
|
||||
var leads = new java.util.LinkedHashMap<String, String>();
|
||||
leads.put("term_other", "some-other-lead");
|
||||
@@ -142,7 +143,7 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void holdsWhenSeveralLeadsAreKnownAndNoneCarriesTheCoordId() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
|
||||
var leads = new java.util.LinkedHashMap<String, String>();
|
||||
leads.put("term_one", "lead-one");
|
||||
leads.put("term_two", "lead-two");
|
||||
@@ -159,7 +160,7 @@ class LeadCoordLoopTest {
|
||||
var channel = new FakeLeadChannel(SELF)
|
||||
.hold(new LeadMessage("m1", PEER, SELF, "first"))
|
||||
.hold(new LeadMessage("m2", PEER, SELF, "second"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
|
||||
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
|
||||
|
||||
loop.tick();
|
||||
@@ -175,10 +176,51 @@ class LeadCoordLoopTest {
|
||||
@Test
|
||||
void anEmptyMailboxNeverTouchesHerdr() {
|
||||
var channel = new FakeLeadChannel(SELF);
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET).agentStatus("idle");
|
||||
|
||||
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
|
||||
|
||||
assertEquals(0, herdr.calls.size(), "an idle fleet must not poll a pane's status every tick");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aLeadWithUnsubmittedTextInItsPromptBoxKeepsTheMessageHeldAndUnacked() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.DRAFTED_PROMPT_CARET).agentStatus("idle");
|
||||
|
||||
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
|
||||
|
||||
assertEquals(0, prompts(herdr).size(),
|
||||
"delivery pastes and submits, so it must not land on a half-typed line");
|
||||
assertEquals(List.of(), channel.acked(), "an undelivered message stays on the broker");
|
||||
assertFalse(channel.peek().isEmpty(), "and is still held");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aMessageHeldForADraftIsDeliveredOnALaterTick() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
|
||||
var herdr = new FakeHerdr().detectionText(FakeHerdr.DRAFTED_PROMPT_CARET).agentStatus("idle");
|
||||
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
|
||||
|
||||
loop.tick();
|
||||
assertEquals(0, prompts(herdr).size());
|
||||
|
||||
herdr.detectionText(FakeHerdr.IDLE_PROMPT_CARET);
|
||||
loop.tick();
|
||||
|
||||
assertEquals(1, prompts(herdr).size(), "the held message lands once the box is empty");
|
||||
assertEquals(List.of("m1"), channel.acked());
|
||||
}
|
||||
|
||||
@Test
|
||||
void anUnreadablePaneKeepsTheMessageHeld() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
|
||||
var herdr = new FakeHerdr().detectionText("garbled ansi noise with no input box").agentStatus("idle");
|
||||
|
||||
loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick();
|
||||
|
||||
assertEquals(0, prompts(herdr).size(), "a pane whose box cannot be found may be holding a draft");
|
||||
assertEquals(List.of(), channel.acked());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
@@ -695,6 +696,8 @@ class LeadHeartbeatLoopTest {
|
||||
private final List<String> sentTexts = new ArrayList<>();
|
||||
private final List<String> promptTargets = new ArrayList<>();
|
||||
private boolean throwOnNextSend = false;
|
||||
/** What {@code agent.read} reports — the loop reads the lead's input box before it nudges. */
|
||||
private String paneTail = FakeHerdr.IDLE_PROMPT_CARET;
|
||||
|
||||
FailableHerdrClient(String lead) {
|
||||
this.lead = lead;
|
||||
@@ -704,6 +707,10 @@ class LeadHeartbeatLoopTest {
|
||||
throwOnNextSend = true;
|
||||
}
|
||||
|
||||
void paneTail(String tail) {
|
||||
this.paneTail = tail;
|
||||
}
|
||||
|
||||
List<String> sentTexts() {
|
||||
return List.copyOf(sentTexts);
|
||||
}
|
||||
@@ -722,6 +729,10 @@ class LeadHeartbeatLoopTest {
|
||||
.put("terminal_id", lead)
|
||||
.put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
return MAPPER.createObjectNode()
|
||||
.set("read", MAPPER.createObjectNode().put("text", paneTail));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
|
||||
if (throwOnNextSend) {
|
||||
@@ -738,4 +749,49 @@ class LeadHeartbeatLoopTest {
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
// ── the operator's own prompt box ────────────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void tickHoldsTheNudgeWhileTheLeadsPromptBoxHoldsUnsubmittedText() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
herdr.paneTail(FakeHerdr.DRAFTED_PROMPT_CARET);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()};
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
|
||||
|
||||
loop.tick(); // opens the idle window
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
|
||||
loop.tick(); // would INJECT, but the operator is mid-sentence
|
||||
assertEquals(0, herdr.sentTexts().size(),
|
||||
"a nudge pastes and submits, so it must not land on a half-typed line");
|
||||
|
||||
herdr.paneTail(FakeHerdr.IDLE_PROMPT_CARET);
|
||||
loop.tick();
|
||||
assertEquals(1, herdr.sentTexts().size(), "the held nudge lands once the box is empty");
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"),
|
||||
"and it still carries the notice the held tick did not spend: " + herdr.sentTexts().get(0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void anUnreadablePaneHoldsTheHeartbeatNudge() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
herdr.paneTail("garbled ansi noise with no input box");
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()};
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
|
||||
|
||||
loop.tick();
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
loop.tick();
|
||||
|
||||
assertEquals(0, herdr.sentTexts().size(),
|
||||
"a pane whose box cannot be found may be holding a draft");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -2376,7 +2376,7 @@ class MessageServiceTest {
|
||||
private PushWiring wireWithPushLoop(int maxReminders, long backoffMs, java.util.function.LongSupplier nowNanos) {
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(T, LEAD);
|
||||
FakeHerdr leadHerdr = new FakeHerdr();
|
||||
FakeHerdr leadHerdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET);
|
||||
AgentControl leadAgents = new AgentControl(leadHerdr);
|
||||
var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor();
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
|
||||
@@ -2413,7 +2413,7 @@ class MessageServiceTest {
|
||||
java.util.function.LongSupplier nowNanos) {
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(T, LEAD);
|
||||
FakeHerdr leadHerdr = new FakeHerdr();
|
||||
FakeHerdr leadHerdr = new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET);
|
||||
AgentControl leadAgents = new AgentControl(leadHerdr);
|
||||
ManualScheduler scheduler = new ManualScheduler();
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
|
||||
@@ -2672,7 +2672,8 @@ class MessageServiceTest {
|
||||
var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor();
|
||||
// A backoff far longer than the test: the schedule is started but no tick ever fires, so
|
||||
// decide() is read directly and nothing here depends on timing.
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(new FakeHerdr()), inbox,
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry,
|
||||
new AgentControl(new FakeHerdr().detectionText(FakeHerdr.IDLE_PROMPT_CARET)), inbox,
|
||||
scheduler, 5, 60_000);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop);
|
||||
try {
|
||||
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.msg;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
@@ -1079,6 +1080,76 @@ class ReplyPushLoopTest {
|
||||
"hitting the ticket reminder cap must count as exhausted");
|
||||
}
|
||||
|
||||
// --- the operator's own prompt box ----------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void aLeadWithUnsubmittedTextInItsPromptBoxIsNotNudged() {
|
||||
var herdr = new PaneTextHerdrClient(FakeHerdr.DRAFTED_PROMPT_CARET);
|
||||
agents = new AgentControl(herdr);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onReplyQueued(WORKER);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0),
|
||||
"a nudge pastes and submits, so an idle lead mid-sentence must not be nudged");
|
||||
assertEquals(0, herdr.promptCount(), "nothing reached the pane");
|
||||
assertFalse(inbox.peek(WORKER).isEmpty(), "and the reply is still waiting to be collected");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theSameLeadIsNudgedOnceItsPromptBoxIsEmpty() {
|
||||
var herdr = new PaneTextHerdrClient(FakeHerdr.DRAFTED_PROMPT_CARET);
|
||||
agents = new AgentControl(herdr);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onReplyQueued(WORKER);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
|
||||
herdr.paneText(FakeHerdr.IDLE_PROMPT_CARET);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
|
||||
"the box emptied, so the held nudge is due");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anUnrecognisablePaneHoldsTheNudge() {
|
||||
var herdr = new PaneTextHerdrClient("garbled ansi noise with no input box");
|
||||
agents = new AgentControl(herdr);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onReplyQueued(WORKER);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0),
|
||||
"a pane whose box cannot be found may be holding a draft");
|
||||
assertEquals(0, herdr.promptCount());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFailedPaneReadHoldsTheNudge() {
|
||||
var herdr = new PaneTextHerdrClient(FakeHerdr.IDLE_PROMPT_CARET).failReads();
|
||||
agents = new AgentControl(herdr);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onReplyQueued(WORKER);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0),
|
||||
"an unreadable box is treated as a draft, never as an empty one");
|
||||
assertEquals(0, herdr.promptCount());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNudgeHeldForADraftIsSentOnALaterTick() throws Exception {
|
||||
var herdr = new PaneTextHerdrClient(FakeHerdr.DRAFTED_PROMPT_CARET);
|
||||
agents = new AgentControl(herdr);
|
||||
|
||||
loop(5, 50).onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
Thread.sleep(300);
|
||||
assertEquals(0, herdr.promptCount(), "every tick holds while the operator is typing");
|
||||
herdr.paneText(FakeHerdr.IDLE_PROMPT_CARET);
|
||||
assertTrue(herdr.sendLatch.await(3, TimeUnit.SECONDS),
|
||||
"the nudge lands on the first tick after the box empties");
|
||||
}
|
||||
|
||||
// --- helpers -------------------------------------------------------------------------------
|
||||
|
||||
private ReplyPushLoop loop() {
|
||||
@@ -1097,6 +1168,15 @@ class ReplyPushLoopTest {
|
||||
return new AgentControl(new FakeHerdrClient(status));
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@code agent.read} frame every fake here returns: a lead settled at an empty input box. The
|
||||
* loop reads the box before it nudges, so a fake that answered nothing would read as a pane it
|
||||
* cannot classify and hold every nudge.
|
||||
*/
|
||||
private static JsonNode emptyPromptBoxRead() {
|
||||
return MAPPER.createObjectNode().set("read", MAPPER.createObjectNode().put("text", FakeHerdr.IDLE_PROMPT_CARET));
|
||||
}
|
||||
|
||||
/** Non-recording (single-threaded) fake — safe for decide() tests. */
|
||||
private static final class FakeHerdrClient implements HerdrClient {
|
||||
private final String agentStatus;
|
||||
@@ -1113,6 +1193,9 @@ class ReplyPushLoopTest {
|
||||
.put("terminal_id", PRIMARY)
|
||||
.put("agent_status", agentStatus));
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
return emptyPromptBoxRead();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@@ -1142,6 +1225,9 @@ class ReplyPushLoopTest {
|
||||
calls.add(Map.entry(method, params));
|
||||
sendLatch.countDown();
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
return emptyPromptBoxRead();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@@ -1179,6 +1265,9 @@ class ReplyPushLoopTest {
|
||||
if ("agent.prompt".equals(method)) {
|
||||
sendLatch.countDown();
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
return emptyPromptBoxRead();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@@ -1216,6 +1305,9 @@ class ReplyPushLoopTest {
|
||||
calls.add(Map.entry(method, params));
|
||||
sendLatch.countDown();
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
return emptyPromptBoxRead();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@@ -1266,6 +1358,9 @@ class ReplyPushLoopTest {
|
||||
promptTargets.add(String.valueOf(p.get("target")));
|
||||
sendLatch.countDown();
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
return emptyPromptBoxRead();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@@ -1316,6 +1411,9 @@ class ReplyPushLoopTest {
|
||||
if ("agent.prompt".equals(method)) {
|
||||
promptTargets.add(String.valueOf(p.get("target")));
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
return emptyPromptBoxRead();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@@ -1363,6 +1461,9 @@ class ReplyPushLoopTest {
|
||||
promptTargets.add(String.valueOf(p.get("target")));
|
||||
sendLatch.countDown();
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
return emptyPromptBoxRead();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@@ -1374,4 +1475,60 @@ class ReplyPushLoopTest {
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Thread-safe fake that always reports {@code idle} and serves a mutable pane tail, so a test can
|
||||
* change what the lead's input box holds between ticks. Records every {@code agent.prompt}.
|
||||
*/
|
||||
private static final class PaneTextHerdrClient implements HerdrClient {
|
||||
private final List<Map.Entry<String, Object>> prompts =
|
||||
Collections.synchronizedList(new ArrayList<>());
|
||||
private volatile String paneText;
|
||||
private volatile boolean failReads = false;
|
||||
volatile CountDownLatch sendLatch = new CountDownLatch(1);
|
||||
|
||||
PaneTextHerdrClient(String paneText) {
|
||||
this.paneText = paneText;
|
||||
}
|
||||
|
||||
void paneText(String text) {
|
||||
this.paneText = text;
|
||||
}
|
||||
|
||||
PaneTextHerdrClient failReads() {
|
||||
this.failReads = true;
|
||||
return this;
|
||||
}
|
||||
|
||||
int promptCount() {
|
||||
return prompts.size();
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.get".equals(method)) {
|
||||
return MAPPER.createObjectNode()
|
||||
.set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", PRIMARY)
|
||||
.put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
if (failReads) {
|
||||
throw new HerdrException("herdr socket read timed out");
|
||||
}
|
||||
return MAPPER.createObjectNode()
|
||||
.set("read", MAPPER.createObjectNode().put("text", paneText));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
prompts.add(Map.entry(method, params));
|
||||
sendLatch.countDown();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -91,8 +91,13 @@ class FleetAppAuthTest {
|
||||
MessageService messages = new MessageService(agents, injector, new Rendezvous());
|
||||
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
|
||||
// term_a is the only herdr-owned pane this fixture's PID can resolve to (FakeHerdr's canned
|
||||
// pane list), and this helper's own contract above says that pane is the worker -- so it
|
||||
// must be recognised as a live spawned member here, the same way a real roster would,
|
||||
// rather than falling through to the observer floor.
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, tokenMode, token,
|
||||
Map::of, new MemberRegistry(null));
|
||||
Map::of, new MemberRegistry(null), t -> "term_a".equals(t) ? MemberRole.DEV : null,
|
||||
Map::of);
|
||||
metrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
|
||||
|
||||
app = new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
|
||||
|
||||
@@ -223,6 +223,37 @@ class FleetAppTest {
|
||||
assertTrue(body.has("detail"), res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void agentsReportsTheAgentsTabLabel() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().withTab("w2", "w2:t7", "trinotes");
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
HttpResponse<String> res = req(port, "GET", "/agents");
|
||||
assertEquals(200, res.statusCode(), res.body());
|
||||
JsonNode agents = mapper.readTree(res.body()).get("agents");
|
||||
assertEquals(1, agents.size());
|
||||
assertEquals("sess-1111", agents.get(0).get("sessionId").asText());
|
||||
assertEquals("trinotes", agents.get(0).get("label").asText());
|
||||
}
|
||||
|
||||
/**
|
||||
* The tab-label scan ({@code workspace.list}/{@code tab.list}) is decoration on top of
|
||||
* {@code workers.list()}'s own agent roster, so its failure must not cost that roster: a row
|
||||
* reports a {@code null} label instead, never the {@code herdr_error} envelope.
|
||||
*/
|
||||
@Test
|
||||
void agentsStillReportsTheRosterWhenTheLabelScanFails() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().workspaceListFailsWith("unavailable");
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
HttpResponse<String> res = req(port, "GET", "/agents");
|
||||
assertEquals(200, res.statusCode(), res.body());
|
||||
JsonNode agents = mapper.readTree(res.body()).get("agents");
|
||||
assertEquals(1, agents.size());
|
||||
assertEquals("sess-1111", agents.get(0).get("sessionId").asText());
|
||||
assertTrue(agents.get(0).get("label").isNull(), "a failed label scan must report a null label, not fail the roster: " + res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnWorkerLandsInOwnTabInWorkerSpaceAndInjectsBaseUrl() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
Reference in New Issue
Block a user