Compare commits
20 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 0ff0d6cb02 | |||
| a2cbf50961 | |||
| 11352d0d06 | |||
| c8e31e1195 | |||
| 24782fe329 | |||
| a49d4dcce9 | |||
| 2aacf07841 | |||
| f2e025cb13 | |||
| 820e6f2620 | |||
| 72cc7560f1 | |||
| 50c7674b98 | |||
| fcce44386b | |||
| f3f067f46f | |||
| 9f4b736fb9 | |||
| 52cd0470b6 | |||
| a3eeace447 | |||
| b696c31756 | |||
| 027d413ce9 | |||
| c88f01ecb8 | |||
| 6022aef612 |
@@ -189,8 +189,9 @@ If the first number has grown past 20, somebody has rolled under the restart pat
|
||||
should be replaced with what they measured.
|
||||
|
||||
**Three separate timeouts bound a roll.** `leadRollover.relaunchReadySeconds` (default 45,
|
||||
`FleetConfig.java:1483`) bounds **each** of two waits that run after the relaunch, so the worst case
|
||||
there is about twice that number, not 45 seconds in total. A third bound gives your old pane 10
|
||||
`FleetConfig.java:1483`) bounds **each** of three waits that run after the relaunch, so the worst
|
||||
case there is about three times that number, not 45 seconds in total. The third wait retries
|
||||
`bootstrapText` while herdr answers `agent_not_ready`. A separate bound gives your old pane 10
|
||||
seconds to die (`LeadRollover.PANE_DEATH_TIMEOUT_SECONDS`).
|
||||
|
||||
`fleet_handover{action: "status", token}` answers with one of these:
|
||||
@@ -204,6 +205,7 @@ seconds to die (`LeadRollover.PANE_DEATH_TIMEOUT_SECONDS`).
|
||||
| `RELAUNCH_FAILED` | launching the fresh pane failed |
|
||||
| `RELAUNCH_NEVER_READY` | the fresh pane never became ready within `relaunchReadySeconds` |
|
||||
| `RELAUNCH_NOT_RECOGNISED` | the fresh terminal never resolved as a lead |
|
||||
| `BOOTSTRAP_NEVER_SENT` | the fresh pane was ready, but herdr refused `bootstrapText` with `agent_not_ready` for the whole bound, so the successor never learned where the handover file is |
|
||||
| `FAILED` | the roll threw; `runRollover`'s catch records this rather than leaving it stuck |
|
||||
|
||||
Only `TURN_NEVER_SETTLED` guarantees your context is intact. The other failures can leave you
|
||||
|
||||
@@ -33,9 +33,11 @@ gate uses. A worker also carries its `sessionId`, `profile`, `worktree` and `bra
|
||||
carries the slot name it was bound to; a collaborator carries its registry name and its own
|
||||
`sessionId`, and **no `leader` key** — a collaborator is a named peer, not a primary. An **observer**
|
||||
carries only its own `sessionId`: a pane the daemon could not place as any of the above, authorized
|
||||
to `READ`/`METRICS`, to `REPLY`/`ASK`/`INBOX` on its own pane, and to `SEND` only to a lead or to a target that
|
||||
resolves as an observer too — never to a collaborator, an architect, or a spawned member, and never a
|
||||
ticket. It finds such a target in `fleet_list`'s `panes` array, which for an observer is filtered to
|
||||
to `READ`/`METRICS`, to `REPLY`/`ASK`/`INBOX` on its own pane, to `SEND` only to a lead or to a target
|
||||
that resolves as an observer too — never to a collaborator, an architect, or a spawned member — and to
|
||||
`TASK_READ` only a ticket its own `fleet_send{wait:false}` created. Nothing wider: not another
|
||||
caller's ticket, and not a session's status, which stays refused. It finds such a target in
|
||||
`fleet_list`'s `panes` array, which for an observer is filtered to
|
||||
exactly what it may send to and reduced to `sessionId`, `label`, `status`, `role` and `deliverable`;
|
||||
a lead's row reads `role: "lead"`.
|
||||
Don't infer what you can ask.
|
||||
@@ -71,20 +73,29 @@ and the sender silently receives nothing. Fail toward the recoverable error.
|
||||
never to a spawned member's terminal, while an observer may send only to a lead or another observer
|
||||
pane;
|
||||
reply/ask/inbox are only-as-itself — any peer may answer, or collect its mail, for its own
|
||||
pane, and for no other. A call outside your role is refused, not queued.
|
||||
pane, and for no other. **A ticket is owned by the caller that created it, not by a role**, so an
|
||||
observer or a collaborator may poll a ticket its own `fleet_send{wait:false}` created and no
|
||||
other — never another caller's, and never a session's status, which stays refused for both. A
|
||||
call outside your role is refused, not queued.
|
||||
4. **Delivery is status-gated: one message per turn.** Don't busy-poll a peer's terminal and don't
|
||||
re-send because a call looks slow — the bridge delivers when the peer is `idle`, `blocked` or
|
||||
`done`. A spawned member must **also** have mounted the bridge MCP: until it has, it is not
|
||||
deliverable, and a send waits on that gate for ~60s and then fails without ever reaching its pane.
|
||||
**A lead's own pane has a second gate: its input box must be empty.** The multiplexer pastes and
|
||||
submits in one step, so a delivery that lands while the operator is typing submits their
|
||||
half-written line. A heartbeat, a ticket nudge and lead-to-lead mail therefore wait until the box
|
||||
is clear, and a pane the daemon cannot read as a box waits too. A direct `fleet_send` to a lead's
|
||||
pane does **not** wait for the box yet (fleetd #793); a lead that runs the fleet mod collects
|
||||
that mail instead of having it pasted, so it is not exposed. Nothing is lost — every one of
|
||||
those paths retries — but a lead that leaves text sitting in its box receives nothing until it
|
||||
clears, and the only sign is one warning in `fleetd.out` after 20 held checks in a row. Delivery
|
||||
to a *member* is not gated this way, because nobody types in a member's pane.
|
||||
**A lead's own pane has a second gate, and it is off by default.** The multiplexer pastes and
|
||||
submits in one step, so a delivery that lands while the operator is typing would submit their
|
||||
half-written line. A heartbeat, a ticket nudge and lead-to-lead mail can therefore wait until the
|
||||
box is clear. That check is the config key `promptBoxGateEnabled`; both an absent key and `false`
|
||||
leave it off, so **right now nothing waits on the box**. It is off because detection read the
|
||||
pane's own autocomplete suggestion as the operator's typing, and so held deliveries that were
|
||||
never at risk — 2231 holds in 11 hours across four panes, measured 2026-10-06. Turn it on only
|
||||
once detection is fixed (#797 made it switchable, #802 is the detection bug). While it is on, a
|
||||
lead that leaves text in its box receives nothing until it clears, and the only sign is one
|
||||
warning in `fleetd.out` after 20 held checks in a row; nothing is lost, because every one of
|
||||
those paths retries. A direct `fleet_send` to a lead's pane never waited for the box either way
|
||||
(fleetd #793); a lead that runs the fleet mod collects that mail instead of having it pasted, so
|
||||
it is not exposed. Delivery to a *member* is not gated this way, because nobody types in a
|
||||
member's pane. Re-measure the switch with
|
||||
`grep -n 'promptBoxGateEnabled' fleetd/fleetd.yaml` — no match means it is still off.
|
||||
**A pane that collects its own mail skips the paste.** A session running the fleet mod calls
|
||||
`fleet_inbox` on a timer. While its last call is under 15s old, the daemon queues that pane's
|
||||
next message for collection instead of pasting it, and the mod hands it to Claude with
|
||||
@@ -184,7 +195,11 @@ below are the procedure — run them in order, every task, not only the big ones
|
||||
**Steps 3 and 4 are separate on purpose** — spawning and sending in one loop is how parallel work
|
||||
silently becomes serial, and it is the most common way this layer is wasted. For the same reason,
|
||||
prefer `wait:false` + `fleet_poll` for anything non-trivial: a blocking `fleet_send` is capped by
|
||||
*your own* MCP client call timeout (~60s), well below the task's real runtime.
|
||||
*your own* MCP client call timeout (~60s), well below the task's real runtime. **A blocking send
|
||||
also creates no ticket**, so when it times out there is no id to poll, and a resend can deliver the
|
||||
same brief twice. One route still recovers the answer: a late turn-completion lands in the target's
|
||||
inbox, and `fleet_poll{target}` drains it. That call is `DRAIN`, so only a lead may make it — the
|
||||
receipt names it to a lead and tells every other caller to ask one.
|
||||
|
||||
**Delegating does not delegate responsibility.** Workers open PRs; you are the gate. Never delegate
|
||||
the merge — and merging on a reviewer's word is delegating it by proxy.
|
||||
@@ -291,19 +306,21 @@ simply complies has thrown away the reason there are two of you.
|
||||
|
||||
`fleet_whoami` answered `collaborator`, so your pane's tab matches a `fleet.collaborators.<name>.tab`
|
||||
entry. You are **not** a member: nothing delegates to you, you have no brief, no worktree and no
|
||||
ticket, and **you owe no `fleet_reply`** — the turn contract above is for a session a lead spawned,
|
||||
and it does not apply to you. Read it only to understand what the members around you are doing.
|
||||
ticket raised against you, and **you owe no `fleet_reply`** — the turn contract above is for a
|
||||
session a lead spawned, and it does not apply to you. Read it only to understand what the members around you are doing.
|
||||
|
||||
What you may do: observe the fleet (`fleet_list`, `fleet_profiles`, `fleet_whoami`), and send to a
|
||||
lead or to another collaborator. What you may not: spawn, stop or drain anything, roll a lead's
|
||||
session, answer a member's `fleet_ask`, poll a ticket, or send to a spawned member's terminal. Each
|
||||
of those is refused at the gate, not queued.
|
||||
What you may do: observe the fleet (`fleet_list`, `fleet_profiles`, `fleet_whoami`), send to a lead
|
||||
or to another collaborator, and poll a ticket your own `fleet_send{wait:false}` created. What you may
|
||||
not: spawn, stop or drain anything, roll a lead's session, answer a member's `fleet_ask`, read
|
||||
another caller's ticket or any session's status, or send to a spawned member's terminal. Each of
|
||||
those is refused at the gate, not queued.
|
||||
|
||||
Two limits worth knowing before you hit them. **You cannot reach a worker** — not even to help one —
|
||||
because a worker belongs to the lead that spawned it, and routing around that would make you a
|
||||
second orchestrator with no plan. Send to the lead instead. And **you cannot read a ticket**, so you
|
||||
cannot collect a delegation's reply: `fleet_poll` refuses you at the role gate, and a ticket also
|
||||
records the terminal that created it, so even a leaked id reads nothing.
|
||||
second orchestrator with no plan. Send to the lead instead. And **you cannot read another caller's
|
||||
ticket** — only one your own `fleet_send{wait:false}` created. A ticket records the terminal that
|
||||
created it, so a leaked id reads nothing, and a lead's delegation stays the lead's: you cannot
|
||||
collect its reply.
|
||||
|
||||
Being named buys you a channel, not authority. Your `fleet_send` to a lead is coordination between
|
||||
peers: the lead owes you no obedience, and you owe it none.
|
||||
|
||||
@@ -174,6 +174,12 @@ bind:
|
||||
# idleSleepGuard:
|
||||
# enabled: false
|
||||
|
||||
# Whether PromptBox checks a lead's input box before a delivery pastes into it. Default off —
|
||||
# omitting this key (or setting it false) leaves every delivery proceeding without reading the
|
||||
# pane. Detection currently misreads the pane's own autocomplete suggestion at the caret as the
|
||||
# operator's unsubmitted typing, so the check holds deliveries that were never actually at risk.
|
||||
# promptBoxGateEnabled: false
|
||||
|
||||
# herdr Unix socket. Omit to use the client default
|
||||
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
|
||||
@@ -421,7 +421,7 @@ final class FleetdAssembly {
|
||||
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
|
||||
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
|
||||
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
|
||||
pushScheduler, maxReminders, backoffMs, metrics);
|
||||
pushScheduler, maxReminders, backoffMs, metrics, cfg.promptBoxGateEnabled());
|
||||
// fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda
|
||||
// above at the real push loop, now that it exists.
|
||||
pushLoopRef.set(pushLoop);
|
||||
@@ -448,7 +448,8 @@ final class FleetdAssembly {
|
||||
Fleetd.leadContextSource(leadContextGauge, router.leadAgents(), leads,
|
||||
Fleetd.leadConfigDirLookup(() -> config.get().profiles(), leaders),
|
||||
Fleetd.leadContextWindowLookup(() -> config.get().profiles(), leaders)),
|
||||
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm);
|
||||
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm,
|
||||
cfg.promptBoxGateEnabled());
|
||||
heartbeat.start();
|
||||
} else {
|
||||
heartbeat = null;
|
||||
@@ -545,7 +546,7 @@ final class FleetdAssembly {
|
||||
if (leadMailbox != null) {
|
||||
var leadCoordScheduler = ports.newScheduler("bridge-leadcoord-");
|
||||
leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler,
|
||||
LEAD_COORD_INTERVAL_MS);
|
||||
LEAD_COORD_INTERVAL_MS, cfg.promptBoxGateEnabled());
|
||||
leadCoordLoop.start(); // FIFTH of the recurring background loops to start (optional).
|
||||
leadCoordSchedulerRef = leadCoordScheduler;
|
||||
} else {
|
||||
|
||||
@@ -86,29 +86,54 @@ public final class Authz {
|
||||
*/
|
||||
public static final Predicate<String> NO_OBSERVER_SEND_TARGET = target -> false;
|
||||
|
||||
/**
|
||||
* The fail-closed classifier for a ticket-scoped {@code TASK_READ} (an observer's or a
|
||||
* collaborator's): answers no for every ticket, so the grant is refused unless a caller
|
||||
* supplies a real one. {@code MessageService#ownsTicket(String, String)} is the real one, so
|
||||
* a ticket this classifier accepts is one that caller's own {@code fleet_send(wait:false)}
|
||||
* actually created.
|
||||
*/
|
||||
public static final Predicate<String> NO_OWNED_TICKET = ticket -> false;
|
||||
|
||||
/**
|
||||
* Convenience form for a caller with no classifier to supply. Fails closed: a collaborator's
|
||||
* or an observer's {@code SEND} is refused, as if no terminal were a configured lead,
|
||||
* collaborator, or observer-reachable target — the same decision
|
||||
* {@link #NO_KNOWN_LEAD_OR_COLLABORATOR} and {@link #NO_OBSERVER_SEND_TARGET} give explicitly.
|
||||
* Every other action's result is identical to the five-argument form's, since none of them
|
||||
* consult either classifier.
|
||||
* or an observer's {@code SEND}, and an observer's or a collaborator's {@code TASK_READ}, are
|
||||
* refused, as if no terminal were a configured lead, collaborator, or observer-reachable
|
||||
* target, and no ticket were the caller's own — the same decision
|
||||
* {@link #NO_KNOWN_LEAD_OR_COLLABORATOR}, {@link #NO_OBSERVER_SEND_TARGET} and
|
||||
* {@link #NO_OWNED_TICKET} give explicitly. Every other action's result is identical to the
|
||||
* full form's, since none of them consult any of the three classifiers.
|
||||
*
|
||||
* <p>Its default classifiers deny every collaborator and every observer, so a caller
|
||||
* enforcing authorization must use the five-argument form instead.
|
||||
* enforcing authorization must use the full form instead.
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession) {
|
||||
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR, NO_OBSERVER_SEND_TARGET);
|
||||
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR,
|
||||
NO_OBSERVER_SEND_TARGET, NO_OWNED_TICKET);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #permits(Principal, Action, String)}, with a real classifier for a collaborator's
|
||||
* {@code SEND}. An observer's {@code SEND} still fails closed ({@link #NO_OBSERVER_SEND_TARGET}) —
|
||||
* a caller enforcing both grants must use the five-argument form.
|
||||
* {@code SEND}. An observer's {@code SEND}, and an observer's or a collaborator's
|
||||
* {@code TASK_READ}, still fail closed — a caller enforcing all three grants must use the
|
||||
* full form.
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession,
|
||||
Predicate<String> knownLeadOrCollaborator) {
|
||||
return permits(caller, action, targetSession, knownLeadOrCollaborator, NO_OBSERVER_SEND_TARGET);
|
||||
return permits(caller, action, targetSession, knownLeadOrCollaborator,
|
||||
NO_OBSERVER_SEND_TARGET, NO_OWNED_TICKET);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #permits(Principal, Action, String, Predicate)}, with a real classifier for an
|
||||
* observer's {@code SEND} too. An observer's or a collaborator's {@code TASK_READ} still
|
||||
* fails closed — a caller enforcing all three grants must use the full form.
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession,
|
||||
Predicate<String> knownLeadOrCollaborator,
|
||||
Predicate<String> observerSendTarget) {
|
||||
return permits(caller, action, targetSession, knownLeadOrCollaborator, observerSendTarget,
|
||||
NO_OWNED_TICKET);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -116,9 +141,10 @@ public final class Authz {
|
||||
*
|
||||
* @param targetSession the session id in the request path; only consulted for the
|
||||
* worker-scoped actions ({@code REPLY}, {@code ASK},
|
||||
* {@code INBOX}), for a collaborator's {@code SEND}, and for an
|
||||
* observer's {@code SEND}, ignored otherwise, may be
|
||||
* {@code null}
|
||||
* {@code INBOX}), for a collaborator's {@code SEND}, for an
|
||||
* observer's {@code SEND}, and — carrying a ticket id instead
|
||||
* of a session id — for an observer's or a collaborator's
|
||||
* {@code TASK_READ}, ignored otherwise, may be {@code null}
|
||||
* @param knownLeadOrCollaborator whether a terminal is a configured lead or collaborator —
|
||||
* consulted only for a collaborator's {@code SEND}, to confine
|
||||
* it to another named peer and never a spawned member's
|
||||
@@ -128,10 +154,16 @@ public final class Authz {
|
||||
* consulted only for an observer's {@code SEND}, to confine it
|
||||
* to a lead or another observer pane and never a collaborator,
|
||||
* an architect, or a spawned member
|
||||
* @param ticketOwnedByCaller whether {@code targetSession} (here, a ticket id) was created
|
||||
* by this same caller — consulted only for an observer's or a
|
||||
* collaborator's {@code TASK_READ}, to confine it to a ticket
|
||||
* its own {@code fleet_send(wait:false)} created, never another
|
||||
* caller's
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession,
|
||||
Predicate<String> knownLeadOrCollaborator,
|
||||
Predicate<String> observerSendTarget) {
|
||||
Predicate<String> observerSendTarget,
|
||||
Predicate<String> ticketOwnedByCaller) {
|
||||
if (caller == null || caller.isAnonymous()) {
|
||||
return false; // authenticated as nothing ⇒ authorized for nothing
|
||||
}
|
||||
@@ -180,11 +212,17 @@ public final class Authz {
|
||||
|| caller.isCollaborator() || caller.isObserver();
|
||||
|
||||
// Ticket polling and session status, open to every role READ is open to except a
|
||||
// collaborator or an observer. MessageService compares a ticket's creator to the
|
||||
// caller on every read as well, so dropping this gate would not expose another
|
||||
// session's reply — it would move the refusal later and widen what a caller that
|
||||
// never orchestrates can probe.
|
||||
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
|
||||
// collaborator or an observer — with one exception: an observer or a collaborator may
|
||||
// poll a ticket its own fleet_send(wait:false) created, confined by
|
||||
// ticketOwnedByCaller. Session status is not ticket-scoped, so that call site
|
||||
// supplies no real classifier here and an observer's or a collaborator's TASK_READ on
|
||||
// it stays refused. MessageService compares a ticket's creator to the caller on every
|
||||
// read as well, so dropping this gate would not expose another session's reply — it
|
||||
// would move the refusal later and widen what a caller that never orchestrates can
|
||||
// probe.
|
||||
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect()
|
||||
|| ((caller.isObserver() || caller.isCollaborator())
|
||||
&& ticketOwnedByCaller.test(targetSession));
|
||||
|
||||
// fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect
|
||||
// holds READ today (CB-548), so "not primary" must mean not-architect here too — this
|
||||
|
||||
@@ -40,8 +40,9 @@ public enum Role {
|
||||
* fleet.collaborators.<name>.tab} registry). Never spawned — identity comes from the
|
||||
* connection, never a request argument, exactly like {@link #WORKER} and {@link #ARCHITECT}.
|
||||
* May {@code SEND} only to a configured lead or collaborator, {@code REPLY}/{@code ASK} only
|
||||
* as its own pane, and {@code READ}/{@code METRICS}; may not {@code SPAWN}/{@code STOP}/
|
||||
* {@code DRAIN}/{@code HANDOVER}, poll a ticket ({@code TASK_READ}), or reach the
|
||||
* as its own pane, {@code READ}/{@code METRICS}, and {@code TASK_READ} only a ticket its own
|
||||
* {@code fleet_send(wait:false)} created; may not {@code SPAWN}/{@code STOP}/{@code DRAIN}/
|
||||
* {@code HANDOVER}, read a session's status or another caller's ticket, or reach the
|
||||
* coordination broker ({@code COORD_SEND}/{@code COORD_READ}).
|
||||
*/
|
||||
COLLABORATOR,
|
||||
@@ -51,10 +52,11 @@ public enum Role {
|
||||
* configured lead, not a bound architect slot, not a configured collaborator tab. Unforgeable
|
||||
* like a worker's — derived from the connection's pane, never from a request argument, and
|
||||
* honoured regardless of auth mode. May {@code READ} and {@code METRICS}, {@code REPLY}/
|
||||
* {@code ASK} only as its own pane, and {@code SEND} only to a target that would itself
|
||||
* resolve as a lead ({@link #PRIMARY}) or as {@code OBSERVER}; may not {@code SPAWN}/
|
||||
* {@code STOP}/{@code DRAIN}/{@code HANDOVER}, poll a ticket ({@code TASK_READ}), or reach the
|
||||
* coordination broker ({@code COORD_SEND}/{@code COORD_READ}).
|
||||
* {@code ASK} only as its own pane, {@code SEND} only to a target that would itself resolve
|
||||
* as a lead ({@link #PRIMARY}) or as {@code OBSERVER}, and {@code TASK_READ} only a ticket its
|
||||
* own {@code fleet_send(wait:false)} created; may not {@code SPAWN}/{@code STOP}/
|
||||
* {@code DRAIN}/{@code HANDOVER}, read a session's status or another caller's ticket, or reach
|
||||
* the coordination broker ({@code COORD_SEND}/{@code COORD_READ}).
|
||||
*/
|
||||
OBSERVER,
|
||||
|
||||
|
||||
@@ -183,8 +183,9 @@ import java.util.function.Supplier;
|
||||
* recounted again for fleetd #362, again after {@code idleSleepGuard:} was added, again after
|
||||
* {@code models:} was added as deferred, again for fleetd #422, which moved {@code models:}
|
||||
* from deferred to hot-excluded once its on/off half was read live everywhere, and again after
|
||||
* {@code leadRollover:} was added (fleetd #480).</strong>
|
||||
* {@code FleetConfig} has 26 top-level record components: 5 cold, 13 deferred, 3 split, 5
|
||||
* {@code leadRollover:} was added (fleetd #480), and again after {@code promptBoxGateEnabled:}
|
||||
* was added as deferred.</strong>
|
||||
* {@code FleetConfig} has 27 top-level record components: 5 cold, 14 deferred, 3 split, 5
|
||||
* hot-excluded. Five of them are named nowhere in this file, and the reason is the same for all
|
||||
* five: {@code placement}, {@code memberCredentials}, {@code memberLoginShell}, {@code models} and
|
||||
* {@code leadRollover} are <strong>hot</strong> and correctly absent — all five are read live off
|
||||
@@ -288,7 +289,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
static final Set<String> DEFERRED_KEYS = Set.of(
|
||||
"guard", "worktreeRoot", "worktreeGroup", "memberSkills", "primary", "configReload",
|
||||
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
|
||||
"quarantineCooldownSeconds", "profiles", "idleSleepGuard");
|
||||
"quarantineCooldownSeconds", "profiles", "idleSleepGuard", "promptBoxGateEnabled");
|
||||
|
||||
private final Path path;
|
||||
private final AtomicReference<FleetConfig> current;
|
||||
@@ -531,6 +532,12 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
if (!Objects.equals(old.idleSleepGuard(), fresh.idleSleepGuard())) {
|
||||
changed.add("idleSleepGuard");
|
||||
}
|
||||
// ReplyPushLoop, LeadCoordLoop and LeadHeartbeatLoop each build their own PromptBox once, at
|
||||
// startup, off this value — a running loop keeps whichever gate setting it was built with
|
||||
// regardless of a later edit here.
|
||||
if (!Objects.equals(old.promptBoxGateEnabled(), fresh.promptBoxGateEnabled())) {
|
||||
changed.add("promptBoxGateEnabled");
|
||||
}
|
||||
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|
||||
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
|
||||
changed.add("spawnReady*");
|
||||
|
||||
@@ -143,6 +143,13 @@ import java.util.regex.PatternSyntaxException;
|
||||
* @param leadRollover opt-in lead rollover (fleetd #480): {@code null} ⇒ off, and no
|
||||
* {@code dev.ltms.fleet.lead.LeadRollover} is constructed at all — an upgraded
|
||||
* daemon never clears a lead's pane on its own initiative. See {@link LeadRollover}.
|
||||
* @param promptBoxGateEnabled whether {@link dev.ltms.fleet.herdr.PromptBox} checks a lead's input
|
||||
* box before a delivery pastes into it. {@code null} (the default) and
|
||||
* {@code false} both leave the check off: a delivery proceeds without reading
|
||||
* the pane, and the box is never classified. Detection currently misreads the
|
||||
* pane's own autocomplete suggestion at the caret as the operator's
|
||||
* unsubmitted typing, so the check holds deliveries that were never actually
|
||||
* at risk; set this to {@code true} only once that is fixed.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record FleetConfig(
|
||||
@@ -171,7 +178,23 @@ public record FleetConfig(
|
||||
String memberSkills,
|
||||
IdleSleepGuard idleSleepGuard,
|
||||
Models models,
|
||||
LeadRollover leadRollover) {
|
||||
LeadRollover leadRollover,
|
||||
Boolean promptBoxGateEnabled) {
|
||||
|
||||
/** Back-compat form before the {@code promptBoxGateEnabled} key was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload, Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup,
|
||||
String memberLoginShell, String memberSkills, IdleSleepGuard idleSleepGuard,
|
||||
Models models, LeadRollover leadRollover) {
|
||||
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup,
|
||||
memberLoginShell, memberSkills, idleSleepGuard, models, leadRollover, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the {@code leadRollover:} block was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
@@ -331,6 +354,7 @@ public record FleetConfig(
|
||||
} else {
|
||||
profiles = Map.of();
|
||||
}
|
||||
promptBoxGateEnabled = promptBoxGateEnabled != null && promptBoxGateEnabled;
|
||||
}
|
||||
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
@@ -1466,7 +1490,7 @@ public record FleetConfig(
|
||||
* CALLING lead's own turn to end (its pane to report {@code IDLE} or
|
||||
* {@code DONE}) before ending that pane's process at all. See the
|
||||
* paragraph above.
|
||||
* @param relaunchReadySeconds default 45 — bound on EACH of two separate waits that run after
|
||||
* @param relaunchReadySeconds default 45 — bound on EACH of three separate waits that run after
|
||||
* the old lead's pane has been torn down and a fresh one launched: first,
|
||||
* for the fresh pane itself to reach a real turn boundary ({@code IDLE} or
|
||||
* {@code DONE}, never merely {@code BLOCKED}) — the safety gate, since
|
||||
@@ -1479,8 +1503,9 @@ public record FleetConfig(
|
||||
* 10s live), so a budget has to clear more than one scan interval to leave
|
||||
* any real margin for the CLI's own boot time; 20 was rejected for exactly
|
||||
* that reason — at a 10s scan interval it only buys two scans. 45 buys
|
||||
* roughly four. Only a timeout on the FIRST wait (the pane never becomes
|
||||
* ready) withholds {@code bootstrapText}.
|
||||
* roughly four; third, to retry sending {@code bootstrapText} while herdr
|
||||
* reports {@code agent_not_ready}. A timeout on the first wait or the third
|
||||
* one withholds {@code bootstrapText}.
|
||||
* @param bootstrapText default a sentence naming the RESOLVED handover path — sent to the
|
||||
* fresh lead's pane once it reaches a real turn boundary after relaunch,
|
||||
* telling the fresh session where to read the handover and carry on. Left
|
||||
@@ -1936,7 +1961,7 @@ public record FleetConfig(
|
||||
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
|
||||
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "memberSkills",
|
||||
"idleSleepGuard", "models", "leadRollover");
|
||||
"idleSleepGuard", "models", "leadRollover", "promptBoxGateEnabled");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static FleetConfig load(Path path) {
|
||||
@@ -2733,10 +2758,13 @@ public record FleetConfig(
|
||||
// own compact constructor defaults the fields of a block that IS present. Defaulting it
|
||||
// here would construct a LeadRollover object (via Fleetd.java's presence gate) for every
|
||||
// config that never mentioned it.
|
||||
// promptBoxGateEnabled is left as-is: the compact constructor already normalized it to a
|
||||
// real true/false, and false (the gate off) is already the safe default — there is no
|
||||
// further defaulting to do here.
|
||||
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
||||
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, memberSkills,
|
||||
idleSleepGuard, models, leadRollover);
|
||||
idleSleepGuard, models, leadRollover, promptBoxGateEnabled);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -26,6 +26,9 @@ import java.util.regex.Pattern;
|
||||
* <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.
|
||||
*
|
||||
* <p>The check itself can be switched off (see the two-arg constructor): a disabled gate clears
|
||||
* every target without reading its pane, logging a hold, or counting a streak.
|
||||
*/
|
||||
public final class PromptBox {
|
||||
|
||||
@@ -74,19 +77,35 @@ public final class PromptBox {
|
||||
}
|
||||
|
||||
private final AgentControl agents;
|
||||
private final boolean gateEnabled;
|
||||
|
||||
/** Consecutive holds per target, so a box that never clears can be warned about once. */
|
||||
private final Map<String, Integer> holdStreaks = new ConcurrentHashMap<>();
|
||||
|
||||
/** Equivalent to {@link #PromptBox(AgentControl, boolean)} with the gate on. */
|
||||
public PromptBox(AgentControl agents) {
|
||||
this(agents, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param gateEnabled {@code false} makes {@link #clearToSubmit} return {@code true} for every
|
||||
* target without reading its pane, logging a hold, or counting a streak —
|
||||
* every delivery proceeds exactly as if the box were always empty.
|
||||
*/
|
||||
public PromptBox(AgentControl agents, boolean gateEnabled) {
|
||||
this.agents = agents;
|
||||
this.gateEnabled = gateEnabled;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* {@code false} means hold and come back; it never means the delivery failed. Always
|
||||
* {@code true} when the gate is disabled.
|
||||
*/
|
||||
public boolean clearToSubmit(String target) {
|
||||
if (!gateEnabled) {
|
||||
return true;
|
||||
}
|
||||
Reading reading = inspect(target);
|
||||
if (reading.state() == State.EMPTY) {
|
||||
holdStreaks.remove(target);
|
||||
|
||||
@@ -232,7 +232,7 @@ public final class LeadRollover {
|
||||
|
||||
/**
|
||||
* What is known about one token, right now — the answer {@link #status} gives. Distinguishes
|
||||
* five terminal outcomes an approved roll can finish with, one in-flight outcome for a roll
|
||||
* terminal outcomes an approved roll can finish with, one in-flight outcome for a roll
|
||||
* that has been approved but has not finished yet, and two answers for a token that names no
|
||||
* active work at all: still pending confirmation, or nothing known about this token at all.
|
||||
*/
|
||||
@@ -256,7 +256,8 @@ public final class LeadRollover {
|
||||
* that is, in fact, actively running. This is not sticky: the deferred continuation
|
||||
* overwrites this same entry with a terminal state ({@link #ROLLED}, {@link
|
||||
* #TURN_NEVER_SETTLED}, {@link #OLD_PANE_NEVER_DIED}, {@link #RELAUNCH_FAILED}, {@link
|
||||
* #RELAUNCH_NEVER_READY}, {@link #RELAUNCH_NOT_RECOGNISED}, or {@link #FAILED}) once it
|
||||
* #RELAUNCH_NEVER_READY}, {@link #RELAUNCH_NOT_RECOGNISED}, {@link #BOOTSTRAP_NEVER_SENT},
|
||||
* or {@link #FAILED}) once it
|
||||
* finishes — including by throwing, which {@link #runRollover}'s catch turns into {@link
|
||||
* #FAILED} instead of leaving this entry stuck forever.
|
||||
*/
|
||||
@@ -305,6 +306,12 @@ public final class LeadRollover {
|
||||
* an operator should check why the tab was not recognised.
|
||||
*/
|
||||
RELAUNCH_NOT_RECOGNISED,
|
||||
/**
|
||||
* A fresh lead was ready, but herdr kept reporting {@code agent_not_ready} while this class
|
||||
* retried {@code bootstrapText} for {@code relaunchReadySeconds}. The fresh session did not
|
||||
* receive its handover instruction.
|
||||
*/
|
||||
BOOTSTRAP_NEVER_SENT,
|
||||
/**
|
||||
* The deferred continuation threw a {@link RuntimeException} and the continuation thread
|
||||
* died with it. Without this state, that throw would leave {@link #outcomes} holding {@link
|
||||
@@ -731,7 +738,19 @@ public final class LeadRollover {
|
||||
|
||||
IdentityResult identityResult = waitUntilRecognisedAsLead(newAgent.terminalId(),
|
||||
cfg.relaunchReadySeconds());
|
||||
agents.send(newAgent.terminalId(), cfg.bootstrapTextFor(p.handoverPath()));
|
||||
BootstrapResult bootstrapResult = sendBootstrapWithRetry(newAgent.terminalId(),
|
||||
cfg.bootstrapTextFor(p.handoverPath()), cfg.relaunchReadySeconds());
|
||||
if (!bootstrapResult.sent()) {
|
||||
log.warn("lead-rollover: bootstrapText was never sent to fresh terminal {} for lead '{}' "
|
||||
+ "after agent_not_ready persisted for {}ms (token={}, configured={}s)",
|
||||
newAgent.terminalId(), leadName, bootstrapResult.elapsedMillis(), p.token(),
|
||||
cfg.relaunchReadySeconds());
|
||||
outcomes.put(p.token(), new RollStatus(RollState.BOOTSTRAP_NEVER_SENT,
|
||||
"fresh terminal " + newAgent.terminalId() + " kept rejecting bootstrapText with "
|
||||
+ "agent_not_ready for relaunchReadySeconds=" + cfg.relaunchReadySeconds()
|
||||
+ "s (measured elapsed=" + bootstrapResult.elapsedMillis() + "ms)"));
|
||||
return;
|
||||
}
|
||||
if (!identityResult.ready()) {
|
||||
log.warn("lead-rollover: fresh terminal {} for lead '{}' is alive and bootstrapped, but "
|
||||
+ "was never recognised as a live lead — an operator should check why "
|
||||
@@ -755,6 +774,30 @@ public final class LeadRollover {
|
||||
+ newAgent.terminalId()));
|
||||
}
|
||||
|
||||
/**
|
||||
* Sends {@code bootstrapText}, retrying only the transient herdr {@code agent_not_ready} refusal
|
||||
* until {@code readySeconds} elapses. All other failures propagate to {@link #runRollover}.
|
||||
*/
|
||||
private BootstrapResult sendBootstrapWithRetry(String terminal, String bootstrapText, int readySeconds) {
|
||||
long startMillis = nowMillis.getAsLong();
|
||||
long deadline = startMillis + TimeUnit.SECONDS.toMillis(readySeconds);
|
||||
while (nowMillis.getAsLong() < deadline) {
|
||||
try {
|
||||
agents.send(terminal, bootstrapText);
|
||||
return new BootstrapResult(true, nowMillis.getAsLong() - startMillis);
|
||||
} catch (HerdrException e) {
|
||||
if (!"agent_not_ready".equals(e.code())) {
|
||||
throw e;
|
||||
}
|
||||
pollSleeper.run();
|
||||
}
|
||||
}
|
||||
return new BootstrapResult(false, nowMillis.getAsLong() - startMillis);
|
||||
}
|
||||
|
||||
/** The result of {@link #sendBootstrapWithRetry}. */
|
||||
private record BootstrapResult(boolean sent, long elapsedMillis) {}
|
||||
|
||||
/** Attempts {@link #captureAgentWithRetry} makes before letting the failure propagate. */
|
||||
static final int CAPTURE_RETRIES = 3;
|
||||
|
||||
|
||||
@@ -516,7 +516,10 @@ public final class FleetMcp {
|
||||
// Answering a worker's fleet_ask (CB-205): resolve its blocked question and
|
||||
// block for the worker's reply as it resumes the same turn. This is the same
|
||||
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
|
||||
return answer(messages, turnId, content, timeoutMs(a), callerOwner);
|
||||
// ANSWER is primary-or-architect, and an architect is refused DRAIN too, so a
|
||||
// timeout receipt must not name fleet_poll{target} to an architect caller.
|
||||
boolean mayDrainPoll = Authz.permits(caller, Authz.Action.DRAIN, target);
|
||||
return answer(messages, turnId, content, timeoutMs(a), callerOwner, mayDrainPoll);
|
||||
}
|
||||
// CB-548: delegator ownership (which lead's reply nudge this worker routes to,
|
||||
// CB-532) is recorded only once the send is ACCEPTED — MessageService has won the
|
||||
@@ -528,12 +531,18 @@ public final class FleetMcp {
|
||||
// name would never match there anyway, but passing null for them keeps the intent
|
||||
// explicit rather than relying on that lookup to filter it out.
|
||||
String delegatorName = caller.isPrimary() ? caller.name() : null;
|
||||
Runnable onAccepted = () ->
|
||||
primaryRegistry.recordDelegation(target, callerTerminal, delegatorName);
|
||||
// A nudge about this target must never tell its recipient to run a call Authz
|
||||
// would refuse it — asked once, here, while the real Principal is still in
|
||||
// scope, rather than re-derived from a bare terminal string later.
|
||||
boolean mayDrainNudge = Authz.permits(caller, Authz.Action.DRAIN, target);
|
||||
boolean mayAnswerNudge = Authz.permits(caller, Authz.Action.ANSWER, target);
|
||||
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, callerTerminal,
|
||||
delegatorName, mayDrainNudge, mayAnswerNudge);
|
||||
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
|
||||
return Boolean.FALSE.equals(a.get("wait"))
|
||||
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), callerOwner);
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), callerOwner,
|
||||
mayDrainNudge);
|
||||
};
|
||||
// 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,
|
||||
@@ -574,10 +583,19 @@ public final class FleetMcp {
|
||||
Map<String, Object> a = req.arguments();
|
||||
String target = str(a, "target");
|
||||
String coordId = str(a, "coordId");
|
||||
String ticket = str(a, "ticket");
|
||||
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
|
||||
Authz.Action action = toolAction("fleet_poll", a);
|
||||
// TASK_READ reads a ticket, not a DRAIN target -- the gate must see the
|
||||
// ticket id, and an observer's or a collaborator's grant is confined to a
|
||||
// ticket it created.
|
||||
String authzTarget = action == Authz.Action.TASK_READ ? ticket : target;
|
||||
Predicate<String> ticketOwnedByCaller = action == Authz.Action.TASK_READ
|
||||
? t -> messages.ownsTicket(t, principal(exchange).ownerKey())
|
||||
: Authz.NO_OWNED_TICKET;
|
||||
McpSchema.CallToolResult denied = deny(exchange, action, authzTarget, ticketOwnedByCaller);
|
||||
if (denied != null) return denied;
|
||||
return poll(messages, leadChannel, str(a, "ticket"), target, coordId,
|
||||
return poll(messages, leadChannel, ticket, target, coordId,
|
||||
principal(exchange).ownerKey());
|
||||
};
|
||||
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
|
||||
@@ -745,6 +763,16 @@ public final class FleetMcp {
|
||||
return denyFor(principal(exchange), action, target);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #deny(McpSyncServerExchange, Authz.Action, String)}, with a real classifier for
|
||||
* an observer's or a collaborator's {@code TASK_READ} — the only call site that can supply
|
||||
* one is a ticket poll, which knows the ticket id and the caller's owner key.
|
||||
*/
|
||||
private McpSchema.CallToolResult deny(McpSyncServerExchange exchange, Authz.Action action,
|
||||
String target, Predicate<String> ticketOwnedByCaller) {
|
||||
return denyFor(principal(exchange), action, target, ticketOwnedByCaller);
|
||||
}
|
||||
|
||||
/**
|
||||
* The policy half of {@link #deny}: everything except pulling the caller out of the MCP
|
||||
* exchange. Kept separate so the authorization decision — the actual control — is unit-testable
|
||||
@@ -758,13 +786,22 @@ public final class FleetMcp {
|
||||
* @return {@code null} when the call may proceed, or the error result to return when it may not
|
||||
*/
|
||||
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target) {
|
||||
return denyFor(caller, action, target, Authz.NO_OWNED_TICKET);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #denyFor(Principal, Authz.Action, String)}, with a real classifier for an
|
||||
* observer's or a collaborator's {@code TASK_READ}.
|
||||
*/
|
||||
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target,
|
||||
Predicate<String> ticketOwnedByCaller) {
|
||||
// The enforcement switch lives HERE rather than in the exchange-facing wrapper: any future
|
||||
// tool that calls this directly must not be able to skip the gate by accident.
|
||||
if (!authorizationEnforced) {
|
||||
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
|
||||
}
|
||||
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator(),
|
||||
callers.observerSendTarget())) {
|
||||
callers.observerSendTarget(), ticketOwnedByCaller)) {
|
||||
if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) {
|
||||
AuditLog.allowed(caller, action, target); // reads would drown the trail
|
||||
}
|
||||
@@ -999,12 +1036,17 @@ public final class FleetMcp {
|
||||
* {@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.
|
||||
*
|
||||
* (CB-548): {@code onAccepted} records delegator ownership the instant the send is accepted, so
|
||||
* a BUSY interloper never claims a turn it did not win. {@code null} disables recording.
|
||||
* <p>(CB-548): {@code onAccepted} records delegator ownership the instant the send is
|
||||
* accepted, so a BUSY interloper never claims a turn it did not win. {@code null} disables
|
||||
* recording.
|
||||
*
|
||||
* @param mayDrainPoll whether this caller's {@code Authz.Action.DRAIN} is granted — a
|
||||
* TIMED_OUT or BUSY receipt must name {@code fleet_poll} only when this
|
||||
* is {@code true}, since that route is refused otherwise
|
||||
*/
|
||||
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
|
||||
Long timeoutMs, Runnable onAccepted, Set<String> profiles,
|
||||
String callerOwner) {
|
||||
String callerOwner, boolean mayDrainPoll) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
@@ -1014,7 +1056,8 @@ public final class FleetMcp {
|
||||
}
|
||||
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
||||
try {
|
||||
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerOwner), timeout);
|
||||
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerOwner), timeout,
|
||||
sessionId, mayDrainPoll);
|
||||
} catch (HerdrException e) {
|
||||
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
|
||||
}
|
||||
@@ -1025,14 +1068,20 @@ public final class FleetMcp {
|
||||
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
|
||||
* it resumes the same turn — surfaced to the primary identically to a normal send.
|
||||
* {@code callerOwner} must match the turn's recorded owner or this is refused.
|
||||
*
|
||||
* @param mayDrainPoll whether this caller's {@code Authz.Action.DRAIN} is granted — see
|
||||
* {@link #send(MessageService, String, String, Long, Runnable, Set, String,
|
||||
* boolean)}'s same parameter. {@code ANSWER} is primary-or-architect, and
|
||||
* an architect is refused {@code DRAIN}, so this cannot default to true.
|
||||
*/
|
||||
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs,
|
||||
String callerOwner) {
|
||||
String callerOwner, boolean mayDrainPoll) {
|
||||
if (isBlank(turnId) || isBlank(content)) {
|
||||
return error("turnId and content are required to answer a worker's question");
|
||||
}
|
||||
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
||||
return formatReply(messages.answer(turnId, content, timeout, callerOwner), timeout);
|
||||
String target = messages.sessionForTurn(turnId);
|
||||
return formatReply(messages.answer(turnId, content, timeout, callerOwner), timeout, target, mayDrainPoll);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1059,8 +1108,15 @@ public final class FleetMcp {
|
||||
};
|
||||
}
|
||||
|
||||
/** Render a {@link MessageService.Reply} as a tool result — shared by {@link #send} and {@link #answer}. */
|
||||
private static McpSchema.CallToolResult formatReply(MessageService.Reply r, long timeout) {
|
||||
/**
|
||||
* Render a {@link MessageService.Reply} as a tool result — shared by {@link #send} and
|
||||
* {@link #answer}. {@code target} is the worker session the send went to (may be {@code null}
|
||||
* if it could not be resolved), named in a TIMED_OUT_* receipt so the caller knows which
|
||||
* inbox to poll — but only when {@code mayDrainPoll} is {@code true}: naming a route the
|
||||
* caller's own {@code Authz.Action.DRAIN} grant would refuse is worse than naming none.
|
||||
*/
|
||||
private static McpSchema.CallToolResult formatReply(MessageService.Reply r, long timeout, String target,
|
||||
boolean mayDrainPoll) {
|
||||
return switch (r.outcome()) {
|
||||
case REPLIED -> text(r.text());
|
||||
// The worker's turn finished but it never called fleet_reply — hand back the scraped
|
||||
@@ -1083,8 +1139,16 @@ public final class FleetMcp {
|
||||
+ "answered (turnId stale)");
|
||||
case NOT_TURN_OWNER -> error("this turn belongs to a different delegation — only the caller "
|
||||
+ "that opened it may answer it");
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
|
||||
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> {
|
||||
String base = "[no reply within " + timeout + "ms — worker "
|
||||
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + ". "
|
||||
+ MessageService.NO_TICKET_NO_RESEND + ". ";
|
||||
yield text(mayDrainPoll
|
||||
? base + "Poll fleet_poll{target=\""
|
||||
+ (target != null ? target : "<the worker session you sent to>")
|
||||
+ "\"} to drain its inbox for a late reply.]"
|
||||
: base + "A late reply cannot be recovered on this channel — ask the lead.]");
|
||||
}
|
||||
// fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call,
|
||||
// so the message may already be sitting in the pane. Do not invite a blind retry the way
|
||||
// the case above does; a resend on this route can double-deliver the same brief.
|
||||
|
||||
@@ -48,8 +48,13 @@ public final class PrimaryRegistry {
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Delegation> leadByTarget = new ConcurrentHashMap<>();
|
||||
|
||||
/** A recorded delegator: the terminal learned from call traffic, and its name, if it has one. */
|
||||
private record Delegation(String terminal, String name) {
|
||||
/**
|
||||
* A recorded delegator: the terminal learned from call traffic, its name if it has one, and
|
||||
* whether a nudge may tell it to {@code fleet_poll(target=…)} or {@code fleet_send(turnId=…)}
|
||||
* — the real authorization decision for this delegator's role, captured once by the caller at
|
||||
* delegation time rather than re-derived from a bare terminal string later.
|
||||
*/
|
||||
private record Delegation(String terminal, String name, boolean mayDrainNudge, boolean mayAnswerNudge) {
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -128,14 +133,28 @@ public final class PrimaryRegistry {
|
||||
|
||||
/**
|
||||
* As {@link #recordDelegation(String, String)}, additionally recording the delegating lead's
|
||||
* name when the caller carries one. See {@link #record(String, String)} for why the name
|
||||
* matters.
|
||||
* name when the caller carries one, and granting a nudge about this target everything a
|
||||
* primary may be told to run. See {@link #record(String, String)} for why the name matters,
|
||||
* and {@link #recordDelegation(String, String, String, boolean, boolean)} for why every other
|
||||
* caller must supply its own, real grant instead of this default.
|
||||
*/
|
||||
public void recordDelegation(String target, String leadTerminal, String leadName) {
|
||||
recordDelegation(target, leadTerminal, leadName, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #recordDelegation(String, String, String)}, with the delegating caller's own
|
||||
* {@code DRAIN} and {@code ANSWER} grants — the real {@code Authz.permits} decision for its
|
||||
* role, asked once by the caller at delegation time, since a terminal string alone cannot be
|
||||
* resolved back to a role here. A nudge about this target must offer only what
|
||||
* {@link #nudgeMayDrainFor(String)}/{@link #nudgeMayAnswerFor(String)} report back as true.
|
||||
*/
|
||||
public void recordDelegation(String target, String leadTerminal, String leadName,
|
||||
boolean mayDrainNudge, boolean mayAnswerNudge) {
|
||||
if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) {
|
||||
return;
|
||||
}
|
||||
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName)));
|
||||
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName), mayDrainNudge, mayAnswerNudge));
|
||||
}
|
||||
|
||||
/** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */
|
||||
@@ -167,6 +186,32 @@ public final class PrimaryRegistry {
|
||||
return currentPrimaryTerminal();
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether a nudge about {@code target} may tell its recipient to {@code fleet_poll(target=…)}
|
||||
* — present on exactly the same condition as {@link #nudgeTargetFor(String)}, carrying the
|
||||
* grant recorded for a delegation, or {@code true} for the singleton primary fallback, which
|
||||
* is always genuinely the primary.
|
||||
*/
|
||||
public Optional<Boolean> nudgeMayDrainFor(String target) {
|
||||
Delegation delegation = target == null ? null : leadByTarget.get(target);
|
||||
if (delegation != null) {
|
||||
return Optional.of(delegation.mayDrainNudge());
|
||||
}
|
||||
return currentPrimaryTerminal().isPresent() ? Optional.of(true) : Optional.empty();
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #nudgeMayDrainFor(String)}, for whether a nudge may tell its recipient to
|
||||
* {@code fleet_send(turnId=…)}.
|
||||
*/
|
||||
public Optional<Boolean> nudgeMayAnswerFor(String target) {
|
||||
Delegation delegation = target == null ? null : leadByTarget.get(target);
|
||||
if (delegation != null) {
|
||||
return Optional.of(delegation.mayAnswerNudge());
|
||||
}
|
||||
return currentPrimaryTerminal().isPresent() ? Optional.of(true) : Optional.empty();
|
||||
}
|
||||
|
||||
/**
|
||||
* The known primary terminal, or empty if not yet learned (and not pinned) — the raw value as
|
||||
* it was recorded, with no attempt to resolve a named lead's current pane. Callers that need a
|
||||
|
||||
@@ -63,6 +63,12 @@ public final class LeadCoordLoop {
|
||||
|
||||
private volatile boolean running;
|
||||
|
||||
/** Equivalent to {@link #LeadCoordLoop(LeadChannel, AgentControl, Supplier, ScheduledExecutorService, long, boolean)} with the prompt-box gate on. */
|
||||
public LeadCoordLoop(LeadChannel channel, AgentControl agents, Supplier<Map<String, String>> leads,
|
||||
ScheduledExecutorService scheduler, long intervalMs) {
|
||||
this(channel, agents, leads, scheduler, intervalMs, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param channel this daemon's own lead mailbox
|
||||
* @param agents herdr control, for the status gate and the pane injection
|
||||
@@ -71,12 +77,13 @@ public final class LeadCoordLoop {
|
||||
* the tab scan after startup becomes reachable without a restart
|
||||
* @param scheduler the loop's own scheduler; the caller owns its shutdown
|
||||
* @param intervalMs how long between ticks
|
||||
* @param promptBoxGateEnabled passed straight to {@link PromptBox#PromptBox(AgentControl, boolean)}
|
||||
*/
|
||||
public LeadCoordLoop(LeadChannel channel, AgentControl agents, Supplier<Map<String, String>> leads,
|
||||
ScheduledExecutorService scheduler, long intervalMs) {
|
||||
ScheduledExecutorService scheduler, long intervalMs, boolean promptBoxGateEnabled) {
|
||||
this.channel = channel;
|
||||
this.agents = agents;
|
||||
this.promptBox = new PromptBox(agents);
|
||||
this.promptBox = new PromptBox(agents, promptBoxGateEnabled);
|
||||
this.leads = leads;
|
||||
this.scheduler = scheduler;
|
||||
this.intervalMs = intervalMs;
|
||||
|
||||
@@ -137,9 +137,21 @@ public final class LeadHeartbeatLoop {
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||
LeadContextSource contextSource, boolean contextHighNudge,
|
||||
boolean requireOperatorConfirm) {
|
||||
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
|
||||
idleAfterNanos, backoffMs, quietNudgeCap, metrics, contextSource, contextHighNudge,
|
||||
requireOperatorConfirm, true);
|
||||
}
|
||||
|
||||
/** As above, plus the prompt-box gate's on/off switch — see {@link PromptBox#PromptBox(AgentControl, boolean)}. */
|
||||
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||
LeadContextSource contextSource, boolean contextHighNudge,
|
||||
boolean requireOperatorConfirm, boolean promptBoxGateEnabled) {
|
||||
this.primaryRegistry = primaryRegistry;
|
||||
this.agents = agents;
|
||||
this.promptBox = new PromptBox(agents);
|
||||
this.promptBox = new PromptBox(agents, promptBoxGateEnabled);
|
||||
this.inbox = inbox;
|
||||
this.roster = roster;
|
||||
this.pushLoop = pushLoop;
|
||||
|
||||
@@ -134,6 +134,14 @@ public final class MessageService {
|
||||
NOT_TURN_OWNER
|
||||
}
|
||||
|
||||
/**
|
||||
* The fact shared by the MCP and REST receipts for {@link Outcome#TIMED_OUT_WORKING},
|
||||
* {@link Outcome#TIMED_OUT_QUEUED} and {@link Outcome#BUSY}: a blocking send is never tracked
|
||||
* by a ticket, so there is no id to poll by, and resending it risks a duplicate delivery.
|
||||
*/
|
||||
public static final String NO_TICKET_NO_RESEND =
|
||||
"this send created no ticket, and a resend can duplicate the delivery";
|
||||
|
||||
/**
|
||||
* @param outcome how the send ended (or paused)
|
||||
* @param text the worker's answer when {@link #completed()} (a structured {@code fleet_reply}
|
||||
@@ -353,6 +361,16 @@ public final class MessageService {
|
||||
* target ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}).
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* The captured waiter a timed-out send left open for a late turn-completion scrape
|
||||
* ({@link #strandLateResolution}), keyed by target so a late {@code fleet_reply} can find and
|
||||
* claim the exact one its own turn belongs to ({@link #reply}). The waiter's own {@code
|
||||
* complete()} call is the discriminator: whichever caller resolves it first — this reply, or
|
||||
* the eventual completion scrape — decides what the other one sees, and only the winner's kind
|
||||
* ever reaches the inbox. An entry removes itself once its waiter resolves.
|
||||
*/
|
||||
private final ConcurrentHashMap<String, CompletableFuture<Rendezvous.Resolution>> lateWaiters =
|
||||
new ConcurrentHashMap<>();
|
||||
private final AtomicLong ticketSeq = new AtomicLong();
|
||||
/**
|
||||
* Minted once per {@code MessageService} instance and folded into every ticket id (see
|
||||
@@ -612,6 +630,14 @@ public final class MessageService {
|
||||
+ "answers, queuing to the inbox instead of guessing", session, candidates.size(),
|
||||
tickets);
|
||||
}
|
||||
// Claim the matching timed-out send's captured waiter, if one is still open, before
|
||||
// publishing this reply. The waiter is this exact turn's own CompletableFuture, so
|
||||
// completing it here means a completion-fallback scrape that resolves the same waiter
|
||||
// afterward finds it already answered and does not publish a second entry for this turn.
|
||||
CompletableFuture<Rendezvous.Resolution> lateWaiter = lateWaiters.get(session);
|
||||
if (lateWaiter != null) {
|
||||
rendezvous.resolveLateReply(lateWaiter, content);
|
||||
}
|
||||
inbox.publish(session, UUID.randomUUID().toString(), content);
|
||||
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
|
||||
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
|
||||
@@ -808,6 +834,7 @@ public final class MessageService {
|
||||
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
|
||||
strandedReplies.remove(target);
|
||||
queuedDeliveries.remove(target);
|
||||
lateWaiters.remove(target);
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
||||
boolean asyncFailed = false;
|
||||
@@ -1050,6 +1077,7 @@ public final class MessageService {
|
||||
queuedDeliveries.put(target, Boolean.TRUE);
|
||||
outcome = Outcome.TIMED_OUT_QUEUED;
|
||||
}
|
||||
strandLateResolution(target, reply);
|
||||
return recorded(new Reply(outcome, null));
|
||||
} catch (ExecutionException e) {
|
||||
Throwable cause = e.getCause();
|
||||
@@ -1067,6 +1095,42 @@ public final class MessageService {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Once a blocking send has given up on {@code waiter} and already reported a TIMED_OUT_*
|
||||
* outcome, the CB-106 completion fallback can still resolve it later — it completes the exact
|
||||
* captured future ({@link Rendezvous#resolveCompletion}), not a session lookup, so
|
||||
* {@link Rendezvous#close} does not stop it. With nobody left awaiting {@code waiter}, route
|
||||
* that resolution to {@code target}'s inbox instead, exactly where a late explicit
|
||||
* {@code fleet_reply} already lands (see {@link #reply}), so {@code fleet_poll{target}} can
|
||||
* recover it. A no-op if {@code waiter} never resolves, or already has by the time this runs.
|
||||
*
|
||||
* <p>Only {@link Rendezvous.Kind#COMPLETION} is routed here. An explicit {@code fleet_reply}
|
||||
* ({@link Rendezvous.Kind#REPLY}) still reaches {@code waiter} through {@link #reply}'s own
|
||||
* session lookup, not through this captured reference, so leaving it out here avoids a second,
|
||||
* racing publish of the same reply.
|
||||
*/
|
||||
private void strandLateResolution(String target, CompletableFuture<Rendezvous.Resolution> waiter) {
|
||||
lateWaiters.put(target, waiter);
|
||||
waiter.whenComplete((resolution, error) -> {
|
||||
lateWaiters.remove(target, waiter);
|
||||
if (resolution != null && resolution.kind() == Rendezvous.Kind.COMPLETION) {
|
||||
try {
|
||||
// inbox.publish reaches a broker and can throw. An exception thrown inside a
|
||||
// whenComplete callback is swallowed into the discarded dependent stage, not
|
||||
// rethrown to any caller — log it here, or a broker failure is invisible and the
|
||||
// late reply is simply lost.
|
||||
inbox.publish(target, UUID.randomUUID().toString(), resolution.text());
|
||||
strandedReplies.put(target, Boolean.TRUE);
|
||||
if (pushLoop != null) {
|
||||
pushLoop.onReplyQueued(target);
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("strandLateResolution: publishing a late turn-completion for {} failed", target, e);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* A worker's mid-turn question (CB-205 reverse rendezvous): surface {@code question} to the
|
||||
* primary by resolving its open blocking {@code fleet_send}, then block this (worker) call until
|
||||
@@ -1167,6 +1231,14 @@ public final class MessageService {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The worker session an open {@code fleet_ask} turn belongs to, or {@code null} if
|
||||
* {@code turnId} is unknown or has already lapsed.
|
||||
*/
|
||||
public String sessionForTurn(String turnId) {
|
||||
return rendezvous.askSession(turnId);
|
||||
}
|
||||
|
||||
/**
|
||||
* The primary's answer to a worker's {@code fleet_ask} (CB-205): resolve the worker's blocked
|
||||
* question identified by {@code turnId}, then — like a fresh {@link #send} — block for the worker's
|
||||
@@ -1492,6 +1564,16 @@ public final class MessageService {
|
||||
|| Objects.equals(callerOwner, task.creatorOwner);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code callerOwner} owns {@code ticket}, with no side effect — unlike
|
||||
* {@link #poll(String, String)}, this never fires the ticket's collection hook. A ticket this
|
||||
* daemon has never heard of is owned by nobody.
|
||||
*/
|
||||
public boolean ownsTicket(String ticket, String callerOwner) {
|
||||
Task task = ticket == null ? null : tasks.get(ticket);
|
||||
return task != null && ownsTicket(task, callerOwner);
|
||||
}
|
||||
|
||||
/**
|
||||
* Test seam only — carries no production behaviour, and nothing in this class calls it;
|
||||
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
|
||||
|
||||
@@ -304,6 +304,20 @@ public final class Rendezvous {
|
||||
return waiter != null && waiter.complete(new Resolution(Kind.BACKEND_EXHAUSTED, reason));
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a specific captured {@code waiter} as a reply — for a {@code fleet_reply} that
|
||||
* arrives after the send which opened {@code waiter} already gave up on it, so the ordinary
|
||||
* session-keyed {@link #resolve} finds no live waiter to match. Like
|
||||
* {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured waiter
|
||||
* (CB-116): a no-op if that waiter was already resolved — first resolution wins, and whichever
|
||||
* one wins is what any later resolution attempt on the same waiter will see.
|
||||
*
|
||||
* @return true if this call resolved the waiter, false if it was null or already resolved
|
||||
*/
|
||||
public boolean resolveLateReply(CompletableFuture<Resolution> waiter, String content) {
|
||||
return waiter != null && waiter.complete(new Resolution(Kind.REPLY, content));
|
||||
}
|
||||
|
||||
private boolean complete(String session, Resolution resolution) {
|
||||
ForwardWaiter waiter = waiters.get(session);
|
||||
return waiter != null && waiter.future().complete(resolution);
|
||||
|
||||
@@ -130,9 +130,16 @@ public final class ReplyPushLoop {
|
||||
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
ScheduledExecutorService scheduler,
|
||||
int maxReminders, long backoffMs, Metrics metrics) {
|
||||
this(primaryRegistry, agents, inbox, scheduler, maxReminders, backoffMs, metrics, true);
|
||||
}
|
||||
|
||||
/** As above, with the prompt-box gate's on/off switch — see {@link PromptBox#PromptBox(AgentControl, boolean)}. */
|
||||
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
ScheduledExecutorService scheduler,
|
||||
int maxReminders, long backoffMs, Metrics metrics, boolean promptBoxGateEnabled) {
|
||||
this.primaryRegistry = primaryRegistry;
|
||||
this.agents = agents;
|
||||
this.promptBox = new PromptBox(agents);
|
||||
this.promptBox = new PromptBox(agents, promptBoxGateEnabled);
|
||||
this.inbox = inbox;
|
||||
this.scheduler = scheduler;
|
||||
this.maxReminders = maxReminders;
|
||||
@@ -178,13 +185,20 @@ public final class ReplyPushLoop {
|
||||
pendingReplies.remove(target, owning);
|
||||
continue;
|
||||
}
|
||||
// Never nudge a caller to run a DRAIN its own role could never perform -- the durable
|
||||
// inbox stays the backstop for it, same as when no nudge target is known at all.
|
||||
if (!owning.mayDrainNudge()) continue;
|
||||
result.add(target);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
/** A pending reply target: which lead to nudge, and how many nudges have named it so far. */
|
||||
private record ReplyEntry(String lead, int nudgeCount) {
|
||||
/**
|
||||
* A pending reply target: which nudge target to notify, how many nudges have named it so
|
||||
* far, and whether that nudge target's own role may actually run the
|
||||
* {@code fleet_poll(target=…)} a reply nudge would tell it to run.
|
||||
*/
|
||||
private record ReplyEntry(String lead, int nudgeCount, boolean mayDrainNudge) {
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -207,11 +221,12 @@ public final class ReplyPushLoop {
|
||||
|
||||
/**
|
||||
* An open question awaiting the lead's answer: which ticket it belongs to, which worker asked,
|
||||
* which lead to nudge, the question text, and how many nudges have named it so far (CB-598 —
|
||||
* tracked per question, not per lead per source).
|
||||
* which lead to nudge, the question text, how many nudges have named it so far (CB-598 —
|
||||
* tracked per question, not per lead per source), and whether that nudge target's own role may
|
||||
* actually run the {@code fleet_send(turnId=…)} a question nudge would tell it to run.
|
||||
*/
|
||||
private record PendingQuestion(String turnId, String ticket, String target, String lead,
|
||||
String question, int nudgeCount) {
|
||||
String question, int nudgeCount, boolean mayAnswerNudge) {
|
||||
}
|
||||
|
||||
private record IncidentLead(String incidentId, String lead) {
|
||||
@@ -229,7 +244,11 @@ public final class ReplyPushLoop {
|
||||
|
||||
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
|
||||
private List<PendingQuestion> pendingQuestionsFor(String lead) {
|
||||
return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList();
|
||||
// Never nudge a caller to run an ANSWER its own role could never perform -- the 55-second
|
||||
// ask window lapses on its own, same as when no nudge target is known at all.
|
||||
return pendingQuestions.values().stream()
|
||||
.filter(q -> lead.equals(q.lead()) && q.mayAnswerNudge())
|
||||
.toList();
|
||||
}
|
||||
|
||||
/** Question turnIds still open for {@code lead} — a plain snapshot for race comparison. */
|
||||
@@ -385,7 +404,7 @@ public final class ReplyPushLoop {
|
||||
boolean hasIncidentWork = !pendingIncidentKeysFor(lead).isEmpty();
|
||||
boolean hasUnmappedTargetWork = !pendingUnmappedTargetKeysFor(lead).isEmpty();
|
||||
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork && !hasUnmappedTargetWork) {
|
||||
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
||||
log.debug("push: nothing pending for nudge target {}, stopping reminder", lead);
|
||||
return Action.STOP;
|
||||
}
|
||||
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
|
||||
@@ -394,7 +413,7 @@ public final class ReplyPushLoop {
|
||||
boolean incidentEligible = hasIncidentWork && incidentReminderCount < maxReminders;
|
||||
boolean unmappedTargetEligible = hasUnmappedTargetWork && unmappedTargetReminderCount < maxReminders;
|
||||
if (!replyEligible && !ticketEligible && !questionEligible && !incidentEligible && !unmappedTargetEligible) {
|
||||
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
||||
log.debug("push: reminder cap ({}) reached for nudge target {} on every source with pending work, stopping",
|
||||
maxReminders, lead);
|
||||
countNudge("exhausted");
|
||||
return Action.STOP;
|
||||
@@ -403,15 +422,15 @@ public final class ReplyPushLoop {
|
||||
try {
|
||||
status = agents.status(lead);
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("push: status check failed for lead {}, will retry", lead, e);
|
||||
log.debug("push: status check failed for nudge target {}, will retry", lead, e);
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
if (!status.injectable()) {
|
||||
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
|
||||
log.debug("push: nudge target {} is {} (not injectable), waiting", lead, status);
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
if (!promptBox.clearToSubmit(lead)) {
|
||||
log.debug("push: lead {} has unsubmitted text in its prompt box, waiting", lead);
|
||||
log.debug("push: nudge target {} has unsubmitted text in its prompt box, waiting", lead);
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
return Action.INJECT;
|
||||
@@ -466,14 +485,14 @@ public final class ReplyPushLoop {
|
||||
if (lead.isEmpty() || isLive(lead.get())) {
|
||||
return lead;
|
||||
}
|
||||
log.debug("push: lead {} delegated to for {} is no longer live, forgetting the stale binding "
|
||||
log.debug("push: nudge target {} delegated to for {} is no longer live, forgetting the stale binding "
|
||||
+ "and falling back", lead.get(), target);
|
||||
primaryRegistry.forgetDelegation(target);
|
||||
Optional<String> fallback = primaryRegistry.nudgeTargetFor(target);
|
||||
if (fallback.isEmpty() || isLive(fallback.get())) {
|
||||
return fallback;
|
||||
}
|
||||
log.debug("push: fallback lead {} for {} is also not live, skipping this tick", fallback.get(), target);
|
||||
log.debug("push: fallback nudge target {} for {} is also not live, skipping this tick", fallback.get(), target);
|
||||
return Optional.empty();
|
||||
}
|
||||
|
||||
@@ -496,9 +515,9 @@ public final class ReplyPushLoop {
|
||||
} catch (RuntimeException e) {
|
||||
boolean gone = e instanceof HerdrException he && "agent_not_found".equals(he.code());
|
||||
if (gone) {
|
||||
log.debug("push: lead {} no longer exists ({})", lead, e.toString());
|
||||
log.debug("push: nudge target {} no longer exists ({})", lead, e.toString());
|
||||
} else {
|
||||
log.debug("push: liveness check for lead {} was inconclusive ({}); treating as live "
|
||||
log.debug("push: liveness check for nudge target {} was inconclusive ({}); treating as live "
|
||||
+ "rather than risk destroying a live binding", lead, e.toString());
|
||||
}
|
||||
return !gone;
|
||||
@@ -517,11 +536,12 @@ public final class ReplyPushLoop {
|
||||
public void onReplyQueued(String target) {
|
||||
var lead = resolveLiveLead(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
|
||||
log.debug("push: no nudge target is known to be waiting on {}, skipping reminder", target);
|
||||
return;
|
||||
}
|
||||
boolean mayDrainNudge = primaryRegistry.nudgeMayDrainFor(target).orElse(false);
|
||||
pendingReplies.compute(target, (t, existing) ->
|
||||
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
|
||||
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount(), mayDrainNudge));
|
||||
startOrCoalesce(lead.get());
|
||||
}
|
||||
|
||||
@@ -545,7 +565,7 @@ public final class ReplyPushLoop {
|
||||
public void onTicketTerminal(String ticket, String target, boolean failed) {
|
||||
var lead = resolveLiveLead(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
|
||||
log.debug("push: no nudge target is known to be waiting on ticket {} (target {}), skipping nudge",
|
||||
ticket, target);
|
||||
return;
|
||||
}
|
||||
@@ -580,11 +600,13 @@ public final class ReplyPushLoop {
|
||||
public void onQuestionOpened(String ticket, String target, String turnId, String question) {
|
||||
var lead = resolveLiveLead(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge",
|
||||
log.debug("push: no nudge target is known to be waiting on {}'s question (turnId {}), skipping nudge",
|
||||
target, turnId);
|
||||
return;
|
||||
}
|
||||
pendingQuestions.put(turnId, new PendingQuestion(turnId, ticket, target, lead.get(), question, 0));
|
||||
boolean mayAnswerNudge = primaryRegistry.nudgeMayAnswerFor(target).orElse(false);
|
||||
pendingQuestions.put(turnId,
|
||||
new PendingQuestion(turnId, ticket, target, lead.get(), question, 0, mayAnswerNudge));
|
||||
startOrCoalesce(lead.get());
|
||||
}
|
||||
|
||||
@@ -608,7 +630,7 @@ public final class ReplyPushLoop {
|
||||
for (String target : targets) {
|
||||
var lead = resolveLiveLead(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.warn("push: backend incident {} has no known lead for target {}", incidentId, target);
|
||||
log.warn("push: backend incident {} has no known nudge target for worker {}", incidentId, target);
|
||||
continue;
|
||||
}
|
||||
targetsByLead.computeIfAbsent(lead.get(), _ -> new ArrayList<>()).add(target);
|
||||
@@ -644,10 +666,10 @@ public final class ReplyPushLoop {
|
||||
/** Start a reminder schedule for {@code lead}, or join the one already running. */
|
||||
private void startOrCoalesce(String lead) {
|
||||
if (activeLeads.putIfAbsent(lead, Boolean.TRUE) != null) {
|
||||
log.debug("push: reminder loop already active for lead {}, work coalesced in", lead);
|
||||
log.debug("push: reminder loop already active for nudge target {}, work coalesced in", lead);
|
||||
return;
|
||||
}
|
||||
log.debug("push: starting reminder loop for lead {}", lead);
|
||||
log.debug("push: starting reminder loop for nudge target {}", lead);
|
||||
scheduleNext(lead);
|
||||
}
|
||||
|
||||
@@ -752,11 +774,11 @@ public final class ReplyPushLoop {
|
||||
|| pendingIncidentKeysFor(lead).stream().anyMatch(i -> !incidentsBefore.contains(i))
|
||||
|| pendingUnmappedTargetKeysFor(lead).stream().anyMatch(i -> !unmappedTargetsBefore.contains(i));
|
||||
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
|
||||
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
|
||||
log.debug("push: new work for nudge target {} raced the reminder loop's stop — restarting", lead);
|
||||
scheduleNext(lead);
|
||||
return;
|
||||
}
|
||||
log.debug("push: reminder loop ended for lead {}", lead);
|
||||
log.debug("push: reminder loop ended for nudge target {}", lead);
|
||||
}
|
||||
|
||||
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
||||
@@ -773,13 +795,13 @@ public final class ReplyPushLoop {
|
||||
List<PendingUnmappedTarget> unmappedTargets = pendingUnmappedTargetsFor(lead);
|
||||
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty() && incidents.isEmpty()
|
||||
&& unmappedTargets.isEmpty()) {
|
||||
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
|
||||
log.debug("push: pending work for nudge target {} drained before the nudge could be sent", lead);
|
||||
return;
|
||||
}
|
||||
String nudge = formatNudge(replyTargets, tickets, questions, incidents, unmappedTargets);
|
||||
try {
|
||||
agents.send(lead, nudge);
|
||||
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
|
||||
log.debug("push: nudge sent to nudge target {} (reply {}/{}, ticket {}/{}, question {}/{}; "
|
||||
+ "{} reply target(s), {} ticket(s), {} question(s))",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||
questionReminderCount + 1, maxReminders,
|
||||
@@ -796,7 +818,7 @@ public final class ReplyPushLoop {
|
||||
}
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
|
||||
log.warn("push: failed to nudge target {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||
questionReminderCount + 1, maxReminders, e.toString());
|
||||
}
|
||||
@@ -814,7 +836,8 @@ public final class ReplyPushLoop {
|
||||
List<PendingQuestion> questions, List<PendingIncident> incidents,
|
||||
List<PendingUnmappedTarget> unmappedTargets) {
|
||||
for (String target : replyTargets) {
|
||||
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
|
||||
pendingReplies.computeIfPresent(target,
|
||||
(t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1, e.mayDrainNudge()));
|
||||
}
|
||||
for (PendingTicket ticket : tickets) {
|
||||
pendingTickets.computeIfPresent(ticket.ticket(),
|
||||
@@ -823,7 +846,7 @@ public final class ReplyPushLoop {
|
||||
for (PendingQuestion question : questions) {
|
||||
pendingQuestions.computeIfPresent(question.turnId(), (id, e) ->
|
||||
new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(),
|
||||
e.nudgeCount() + 1));
|
||||
e.nudgeCount() + 1, e.mayAnswerNudge()));
|
||||
}
|
||||
for (PendingIncident incident : incidents) {
|
||||
pendingIncidents.computeIfPresent(incident.key(), (id, e) -> new PendingIncident(e.key(),
|
||||
|
||||
@@ -293,6 +293,20 @@ public final class FleetApp {
|
||||
return Authz.permits(caller, action, target, knownLeadOrCollaborator, observerSendTarget);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #permitsFor(Principal, Authz.Action, String, Predicate, Predicate)}, also
|
||||
* threading the classifier an observer's or a collaborator's {@code TASK_READ} is checked
|
||||
* against; pass a real {@code MessageService#ownsTicket(String, String)}-backed predicate to
|
||||
* exercise the real production gate for a ticket poll, as {@link #allow} does.
|
||||
*/
|
||||
static boolean permitsFor(Principal caller, Authz.Action action, String target,
|
||||
Predicate<String> knownLeadOrCollaborator,
|
||||
Predicate<String> observerSendTarget,
|
||||
Predicate<String> ticketOwnedByCaller) {
|
||||
return Authz.permits(caller, action, target, knownLeadOrCollaborator, observerSendTarget,
|
||||
ticketOwnedByCaller);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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}.
|
||||
@@ -302,11 +316,22 @@ public final class FleetApp {
|
||||
* yours" (a worker reaching for another worker's session, or for orchestration).
|
||||
*/
|
||||
private boolean allow(Context ctx, Authz.Action action, String target) {
|
||||
return allow(ctx, action, target, Authz.NO_OWNED_TICKET);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #allow(Context, Authz.Action, String)}, with a real classifier for an observer's
|
||||
* or a collaborator's {@code TASK_READ} — the only call site that can supply one is a ticket
|
||||
* poll, which knows the ticket id and the caller's owner key.
|
||||
*/
|
||||
private boolean allow(Context ctx, Authz.Action action, String target,
|
||||
Predicate<String> ticketOwnedByCaller) {
|
||||
if (auth == null) {
|
||||
return true; // legacy: authorization not enforced
|
||||
}
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.observerSendTarget())) {
|
||||
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.observerSendTarget(),
|
||||
ticketOwnedByCaller)) {
|
||||
if (action != Authz.Action.READ && action != Authz.Action.METRICS
|
||||
&& action != Authz.Action.TASK_READ) {
|
||||
AuditLog.allowed(caller, action, target); // reads would drown the trail
|
||||
@@ -717,9 +742,15 @@ public final class FleetApp {
|
||||
content = FleetMcp.attributeIfObserver(caller, content);
|
||||
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
|
||||
|
||||
// Whether THIS caller's own Authz.Action.DRAIN is granted — a TIMED_OUT_*/BUSY receipt
|
||||
// must not name GET /sessions/{id}/replies when that same caller would be refused it.
|
||||
// Mirrors allow()'s own legacy bypass: with no CallerResolver configured, that route is
|
||||
// not gated at all, so every caller may reach it.
|
||||
boolean mayDrainPoll = auth == null || Authz.permits(caller, Authz.Action.DRAIN, id);
|
||||
|
||||
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
|
||||
if (turnId != null && !turnId.isBlank()) {
|
||||
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerOwner), timeout);
|
||||
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerOwner), timeout, mayDrainPoll);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -731,7 +762,7 @@ public final class FleetApp {
|
||||
}
|
||||
|
||||
try {
|
||||
writeReply(ctx, id, messages.send(id, content, timeout, callerOwner), timeout);
|
||||
writeReply(ctx, id, messages.send(id, content, timeout, callerOwner), timeout, mayDrainPoll);
|
||||
} catch (HerdrException e) {
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
@@ -741,8 +772,13 @@ public final class FleetApp {
|
||||
* Render a {@link MessageService.Reply} onto the response — shared by a normal send and a
|
||||
* fleet_ask answer. A structured/scraped completion is 200; a worker's mid-turn question a 202
|
||||
* (with its {@code turnId}); a stale answer a 409; every other non-terminal outcome a typed 202.
|
||||
*
|
||||
* @param mayDrainPoll whether this caller's {@code Authz.Action.DRAIN} is granted — a
|
||||
* TIMED_OUT or BUSY {@code detail} names {@code GET /sessions/{id}/replies}
|
||||
* only when this is {@code true}, since that route shares the same grant.
|
||||
*/
|
||||
private void writeReply(Context ctx, String id, MessageService.Reply reply, long timeout) {
|
||||
private void writeReply(Context ctx, String id, MessageService.Reply reply, long timeout,
|
||||
boolean mayDrainPoll) {
|
||||
switch (reply.outcome()) {
|
||||
case QUESTION -> ctx.status(202).json(Map.of(
|
||||
"sessionId", id, "status", "question",
|
||||
@@ -791,7 +827,10 @@ public final class FleetApp {
|
||||
|| reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED)
|
||||
&& reply.text() != null
|
||||
? reply.text()
|
||||
: "no reply within " + timeout + "ms; poll status or retry"));
|
||||
: "no reply within " + timeout + "ms; " + MessageService.NO_TICKET_NO_RESEND + "; "
|
||||
+ (mayDrainPoll
|
||||
? "GET /sessions/" + id + "/replies drains its inbox for a late reply"
|
||||
: "a late reply cannot be recovered on this channel — ask the lead")));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -928,11 +967,15 @@ public final class FleetApp {
|
||||
|
||||
/** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */
|
||||
private void taskStatus(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
|
||||
String ticket = ctx.pathParam("ticket");
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
String callerOwner = caller == null ? null : caller.ownerKey();
|
||||
// The gate must see the ticket id, not a literal null -- an observer's or a
|
||||
// collaborator's grant is confined to a ticket its own caller created.
|
||||
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), ticket, t -> messages.ownsTicket(t, callerOwner))) {
|
||||
return;
|
||||
}
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.ownerKey());
|
||||
MessageService.TaskView v = messages.poll(ticket, callerOwner);
|
||||
if (v == null) {
|
||||
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
|
||||
return;
|
||||
|
||||
@@ -259,18 +259,39 @@ class AuthzTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* Every action denied to a collaborator, asserted denied even when the classifier would
|
||||
* accept any target — proving none of these is actually gated on the classifier at all.
|
||||
* Every action denied to a collaborator, asserted denied even when the
|
||||
* {@code knownLeadOrCollaborator} classifier would accept any target — proving none of these
|
||||
* is actually gated on that classifier at all. {@code TASK_READ} is excluded here and given
|
||||
* its own matrix below, since — unlike every action in this loop — its grant is conditional
|
||||
* on the separate ticket-ownership classifier, not on this one.
|
||||
*/
|
||||
@Test
|
||||
void aCollaboratorIsDeniedLifecycleCoordinationAndTicketPolling() {
|
||||
void aCollaboratorIsDeniedLifecycleAndCoordination() {
|
||||
for (Authz.Action a : new Authz.Action[]{SPAWN, STOP, DRAIN, HANDOVER, ANSWER, COORD_SEND,
|
||||
COORD_READ, TASK_READ}) {
|
||||
COORD_READ}) {
|
||||
assertFalse(Authz.permits(COLLABORATOR, a, "term_lead", target -> true),
|
||||
"a collaborator must not " + a + " even when the classifier accepts every target");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code TASK_READ} is conditional for a collaborator too (fleetd #804), on the same
|
||||
* ticket-ownership classifier an observer's grant uses: it may poll a ticket its own
|
||||
* {@code fleet_send(wait:false)} created, and nothing else.
|
||||
*/
|
||||
@Test
|
||||
void aCollaboratorMayTaskReadOnlyATicketTheClassifierSaysItOwns() {
|
||||
assertTrue(Authz.permits(COLLABORATOR, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
|
||||
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
|
||||
"the classifier accepting the ticket as the caller's own must grant TASK_READ");
|
||||
assertFalse(Authz.permits(COLLABORATOR, TASK_READ, "someone_elses_ticket",
|
||||
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, Authz.NO_OBSERVER_SEND_TARGET, target -> false),
|
||||
"the classifier refusing the ticket must deny TASK_READ");
|
||||
assertFalse(Authz.permits(COLLABORATOR, TASK_READ, "own-ticket"),
|
||||
"the three-argument convenience form fails closed, so TASK_READ is refused without "
|
||||
+ "a real classifier");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aCollaboratorIsNotCountedAsPrimaryWorkerOrArchitect() {
|
||||
assertFalse(COLLABORATOR.isPrimary());
|
||||
@@ -310,17 +331,17 @@ class AuthzTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* Every action beyond READ/METRICS/REPLY/ASK/INBOX/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. The exempt set is the
|
||||
* three only-as-itself actions plus the two open reads. {@code SEND} is excluded here and
|
||||
* given its own matrix below, since — unlike every action in this loop — its grant is
|
||||
* conditional on the target, not fixed.
|
||||
* Every action beyond READ/METRICS/REPLY/ASK/INBOX/SEND, asserted denied for an observer.
|
||||
* The exempt set is the three only-as-itself actions plus the two open reads. {@code SEND}
|
||||
* and {@code TASK_READ} are excluded here and given their own matrices below, since — unlike
|
||||
* every action in this loop — each one's grant is conditional on more than the caller's role
|
||||
* alone.
|
||||
*/
|
||||
@Test
|
||||
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxAndSend() {
|
||||
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxSendAndTaskRead() {
|
||||
for (Authz.Action a : Authz.Action.values()) {
|
||||
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND) {
|
||||
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND
|
||||
|| a == TASK_READ) {
|
||||
continue;
|
||||
}
|
||||
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
|
||||
@@ -328,6 +349,49 @@ class AuthzTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code TASK_READ} is conditional for an observer, not fixed: it may poll a ticket its own
|
||||
* {@code fleet_send(wait:false)} created, confined by {@code ticketOwnedByCaller}, and nothing
|
||||
* else — a session's status is not ticket-scoped, so no call site can ever supply a predicate
|
||||
* that grants it, and the 3-/4-/5-argument convenience forms stay fail-closed for it too.
|
||||
*/
|
||||
@Test
|
||||
void anObserverMayTaskReadOnlyATicketTheClassifierSaysItOwns() {
|
||||
assertTrue(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
|
||||
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
|
||||
"the classifier accepting the ticket as the caller's own must grant TASK_READ");
|
||||
assertFalse(Authz.permits(OBSERVER, TASK_READ, "someone_elses_ticket",
|
||||
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, Authz.NO_OBSERVER_SEND_TARGET, target -> false),
|
||||
"the classifier refusing the ticket must deny TASK_READ");
|
||||
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket"),
|
||||
"the three-argument convenience form fails closed, so TASK_READ is refused without "
|
||||
+ "a real classifier");
|
||||
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR),
|
||||
"the four-argument form still fails closed for TASK_READ");
|
||||
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
|
||||
Authz.NO_OBSERVER_SEND_TARGET),
|
||||
"the five-argument form still fails closed for TASK_READ, since it supplies no "
|
||||
+ "ticketOwnedByCaller classifier either");
|
||||
}
|
||||
|
||||
/**
|
||||
* Control for the test above: every other action's result for an observer does not move when
|
||||
* the ticket-ownership classifier does. Only {@code TASK_READ} is wired to it.
|
||||
*/
|
||||
@Test
|
||||
void theTicketOwnershipClassifierMovesOnlyTaskReadForAnObserver() {
|
||||
for (Authz.Action a : Authz.Action.values()) {
|
||||
if (a == TASK_READ) {
|
||||
continue;
|
||||
}
|
||||
assertEquals(
|
||||
Authz.permits(OBSERVER, a, "term_observer"),
|
||||
Authz.permits(OBSERVER, a, "term_observer", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
|
||||
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
|
||||
a + " must not depend on the ticket-ownership classifier at all");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anObserverIsNotCountedAsAnyOtherRole() {
|
||||
assertFalse(OBSERVER.isPrimary());
|
||||
|
||||
@@ -114,6 +114,7 @@ class ConfigRefTopLevelReportingCoverageTest {
|
||||
// fleetd #480: leadRollover joins placement/memberCredentials/memberLoginShell as
|
||||
// hot-excluded — never compared by any changed*Keys method, so it stays null like them.
|
||||
v.put("leadRollover", null);
|
||||
v.put("promptBoxGateEnabled", true);
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
@@ -159,6 +160,7 @@ class ConfigRefTopLevelReportingCoverageTest {
|
||||
v.put("models", new FleetConfig.Models(
|
||||
List.of(new FleetConfig.Models.ModelEntry("model-b"))));
|
||||
v.put("leadRollover", null);
|
||||
v.put("promptBoxGateEnabled", false);
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
+1
@@ -104,6 +104,7 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
|
||||
// here proves it, rather than leaving it null and proving nothing.
|
||||
v.put("leadRollover", new FleetConfig.LeadRollover(
|
||||
"/handover/guard.md", true, 3600, 20, 45, "read the handover file"));
|
||||
v.put("promptBoxGateEnabled", true);
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
@@ -177,4 +177,24 @@ class PromptBoxTest {
|
||||
herdr.detectionText(EMPTY);
|
||||
assertTrue(box.clearToSubmit("term_a"));
|
||||
}
|
||||
|
||||
// --- the gate's on/off switch ---------------------------------------------
|
||||
|
||||
@Test
|
||||
void aDisabledGateClearsADraftedBoxWithoutReadingThePane() {
|
||||
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
|
||||
|
||||
assertTrue(new PromptBox(new AgentControl(herdr), false).clearToSubmit("term_a"),
|
||||
"a disabled gate must clear every target, draft or not");
|
||||
assertFalse(herdr.called("agent.read"), "a disabled gate must never read the pane");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anEnabledGateBehavesExactlyLikeTheSingleArgConstructor() {
|
||||
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
|
||||
|
||||
assertFalse(new PromptBox(new AgentControl(herdr), true).clearToSubmit("term_a"),
|
||||
"gateEnabled=true must hold a draft exactly like the default constructor");
|
||||
assertTrue(herdr.called("agent.read"), "an enabled gate still reads the pane");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,10 +2,12 @@ package dev.ltms.fleet.lead;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
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.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
@@ -20,6 +22,7 @@ import java.util.HashMap;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -1323,6 +1326,66 @@ class LeadRolloverTest {
|
||||
+ "so an operator reading status() has something to act on: " + status.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BOOTSTRAP 1] one agent_not_ready bootstrap refusal is retried and sends the "
|
||||
+ "handover instruction exactly once")
|
||||
void transientAgentNotReadyRetriesBootstrapAndRolls() throws IOException {
|
||||
FakeHerdr fake = herdrReadyForAFullRoll();
|
||||
AtomicInteger refusedSends = new AtomicInteger();
|
||||
HerdrClient transientRefusal = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.prompt".equals(method) && refusedSends.getAndIncrement() == 0) {
|
||||
throw new HerdrException("herdr error [agent_not_ready]: agent.prompt failed",
|
||||
"agent_not_ready", null);
|
||||
}
|
||||
return fake.call(method, params);
|
||||
}
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
};
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
AgentControl agents = new AgentControl(transientRefusal);
|
||||
WorkspaceControl spaces = new WorkspaceControl(transientRefusal);
|
||||
LeadRollover rollover = new LeadRollover(agents, spaces,
|
||||
new LeadLauncher(agents, spaces, fleetConfigWithRelaunchableLead()),
|
||||
() -> cfg(handover.toString()), _ -> null, _ -> LEAD_NAME,
|
||||
() -> Map.of(NEW_TERMINAL, LEAD_NAME), fixedClock(clock), () -> { }, Runnable::run);
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
assertTrue(rollover.confirm(LEAD, pending.token(), true).accepted());
|
||||
|
||||
assertEquals(2, refusedSends.get(), "one agent_not_ready refusal must be followed by one retry");
|
||||
assertEquals(1, promptCallCount(fake), "only the successful retry reaches herdr delivery");
|
||||
assertEquals(LeadRollover.RollState.ROLLED, rollover.status(pending.token()).state());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BOOTSTRAP 2] persistent agent_not_ready records BOOTSTRAP_NEVER_SENT and "
|
||||
+ "releases the single-flight claim")
|
||||
void persistentAgentNotReadyRecordsBootstrapNeverSentAndReleasesClaim() throws IOException {
|
||||
FakeHerdr fake = herdrReadyForAFullRoll();
|
||||
fake.agentSendFailsWith("agent_not_ready");
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRolloverForAFullRoll(fake, cfg(handover.toString()),
|
||||
() -> clock.addAndGet(1_000), Map.of(NEW_TERMINAL, LEAD_NAME));
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
assertTrue(rollover.confirm(LEAD, pending.token(), true).accepted());
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.BOOTSTRAP_NEVER_SENT, status.state());
|
||||
assertNotEquals(LeadRollover.RollState.ROLLED, status.state());
|
||||
assertNotEquals(LeadRollover.RollState.FAILED, status.state());
|
||||
|
||||
LeadRollover.PendingRollover retry = rollover.open(LEAD, "retry after bootstrap timeout");
|
||||
assertTrue(rollover.confirm(LEAD, retry.token(), true).accepted(),
|
||||
"the terminal outcome must release the single-flight claim");
|
||||
}
|
||||
|
||||
// ---- fleetd #726 unit 3: confirm() single-flights one roll at a time per lead ---------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -56,6 +56,7 @@ class FleetMcpAuthzTest {
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private Metrics metrics;
|
||||
private FleetMcp mcp;
|
||||
private MessageService messages;
|
||||
|
||||
@AfterEach
|
||||
void close() {
|
||||
@@ -82,7 +83,7 @@ class FleetMcpAuthzTest {
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
_ -> "tok");
|
||||
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
|
||||
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(),
|
||||
messages = new MessageService(agents, new Injector(agents), new Rendezvous(),
|
||||
new InMemoryReplyInbox());
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
|
||||
metrics = FleetMetrics.create(sessions, new InMemoryReplyInbox());
|
||||
@@ -346,6 +347,73 @@ class FleetMcpAuthzTest {
|
||||
assertNotNull(m.denyFor(OBSERVER, Authz.Action.COORD_SEND, null));
|
||||
}
|
||||
|
||||
/**
|
||||
* An observer may read a ticket its own {@code fleet_send(wait:false)} created, never one
|
||||
* another caller created — exercised against the real {@link MessageService#ownsTicket}, not
|
||||
* a stand-in classifier.
|
||||
*/
|
||||
@Test
|
||||
void anObserverMayTaskReadATicketItCreatedButNotOneAnotherCallerCreated() {
|
||||
FleetMcp m = mcp(true);
|
||||
String ownTicket = messages.sendAsync("term_a", "do it", null, OBSERVER);
|
||||
String othersTicket = messages.sendAsync("term_a", "do it too", null, WORKER_A);
|
||||
|
||||
assertNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, ownTicket,
|
||||
t -> messages.ownsTicket(t, OBSERVER.ownerKey())),
|
||||
"an observer must read a ticket its own send created");
|
||||
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, othersTicket,
|
||||
t -> messages.ownsTicket(t, OBSERVER.ownerKey())),
|
||||
"an observer must not read a ticket a different caller created");
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #anObserverMayTaskReadATicketItCreatedButNotOneAnotherCallerCreated} (fleetd
|
||||
* #804): a collaborator's {@code fleet_send(wait:false)} creates a ticket too, and the only
|
||||
* caller that can ever read it back is the collaborator that created it.
|
||||
*/
|
||||
@Test
|
||||
void aCollaboratorMayTaskReadATicketItCreatedButNotOneAnotherCallerCreated() {
|
||||
FleetMcp m = mcp(true);
|
||||
String ownTicket = messages.sendAsync("term_a", "do it", null, COLLABORATOR);
|
||||
String othersTicket = messages.sendAsync("term_a", "do it too", null, WORKER_A);
|
||||
|
||||
assertNull(m.denyFor(COLLABORATOR, Authz.Action.TASK_READ, ownTicket,
|
||||
t -> messages.ownsTicket(t, COLLABORATOR.ownerKey())),
|
||||
"a collaborator must read a ticket its own send created");
|
||||
assertNotNull(m.denyFor(COLLABORATOR, Authz.Action.TASK_READ, othersTicket,
|
||||
t -> messages.ownsTicket(t, COLLABORATOR.ownerKey())),
|
||||
"a collaborator must not read a ticket a different caller created");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #803: {@code fleet_poll{}} with no {@code ticket} argument resolves to
|
||||
* {@code TASK_READ} with a {@code null} ticket. Before the fix, the real {@code
|
||||
* ownsTicket(String, String)} classifier reached {@code tasks.get(null)} on a
|
||||
* {@code ConcurrentHashMap} and threw a {@code NullPointerException} instead of refusing —
|
||||
* the gate never got to write its {@code AuditLog} refusal, the caller saw an internal error.
|
||||
* Asserts the refusal itself, not merely that nothing throws: a test that only catches the
|
||||
* absence of a throw would also pass if this path quietly started granting the call.
|
||||
*/
|
||||
@Test
|
||||
void anObserverPollingWithNoTicketIsRefusedRatherThanThrowing() {
|
||||
FleetMcp m = mcp(true);
|
||||
McpSchema.CallToolResult denied = m.denyFor(OBSERVER, Authz.Action.TASK_READ, null,
|
||||
t -> messages.ownsTicket(t, OBSERVER.ownerKey()));
|
||||
assertNotNull(denied, "a missing ticket must be refused, not treated as owned");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code fleet_status} stays refused for an observer even once {@code TASK_READ} opens: it is
|
||||
* not ticket-scoped, so no call site supplies a real {@code ticketOwnedByCaller} classifier,
|
||||
* and the default 3-argument {@code denyFor} fails closed.
|
||||
*/
|
||||
@Test
|
||||
void anObserverMayNotReadSessionStatusEvenAfterTaskReadOpensForTickets() {
|
||||
FleetMcp m = mcp(true);
|
||||
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, "term_a"),
|
||||
"a session id is not a ticket the caller could ever own");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theLegacyConstructorLeavesTheGateOpen() {
|
||||
// The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced.
|
||||
|
||||
@@ -78,7 +78,7 @@ class FleetMcpTest {
|
||||
|
||||
private void assertSendRoundTrips(String target, Set<String> profiles) throws Exception {
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null));
|
||||
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null, true));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
@@ -104,7 +104,7 @@ class FleetMcpTest {
|
||||
void sendThenReplyRoundTrips() throws Exception {
|
||||
// fleet_send blocks; fleet_reply resolves it with the worker's structured answer.
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null));
|
||||
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null, true));
|
||||
|
||||
// Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now
|
||||
// queues in the inbox if no waiter is open, which would break the round-trip).
|
||||
@@ -182,7 +182,7 @@ class FleetMcpTest {
|
||||
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null, true));
|
||||
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
|
||||
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
@@ -280,7 +280,7 @@ class FleetMcpTest {
|
||||
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null));
|
||||
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null, true));
|
||||
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
@@ -298,7 +298,7 @@ class FleetMcpTest {
|
||||
@Test
|
||||
void aDifferentCallersMcpAnswerIsRefusedForABlockingSendButTheRealOwnerSucceeds() throws Exception {
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner"));
|
||||
() -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner", true));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
@@ -313,12 +313,12 @@ class FleetMcpTest {
|
||||
String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length());
|
||||
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
|
||||
|
||||
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
|
||||
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker", true);
|
||||
assertTrue(hijacked.isError(), "a caller that did not open this turn must get an error, not an answer");
|
||||
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner", true));
|
||||
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
@@ -357,14 +357,14 @@ class FleetMcpTest {
|
||||
String turnId = asking.turnId();
|
||||
|
||||
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L,
|
||||
"worker:term_attacker");
|
||||
"worker:term_attacker", true);
|
||||
assertTrue(hijacked.isError(), "a caller that did not create this delegation must get an error");
|
||||
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
|
||||
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
|
||||
"a refused answer must not advance the async ticket's phase");
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey()));
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey(), true));
|
||||
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
@@ -383,11 +383,71 @@ class FleetMcpTest {
|
||||
|
||||
@Test
|
||||
void sendTimesOutWithAWorkingNote() {
|
||||
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null);
|
||||
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null, true);
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
|
||||
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #801: a blocking send creates no ticket, ever, so the TIMED_OUT_QUEUED receipt must
|
||||
* name the one route that recovers a late answer — draining the target's inbox — and must not
|
||||
* invite a resend, which can duplicate the delivery.
|
||||
*/
|
||||
@Test
|
||||
void sendTimesOutQueuedNamesFleetPollTargetAndNotTheBareWordRetry() {
|
||||
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null, true);
|
||||
String text = textOf(res);
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
|
||||
assertTrue(text.contains("fleet_poll{target=\"term_a\"}"), "got: " + text);
|
||||
assertFalse(text.contains("retry"), "got: " + text);
|
||||
}
|
||||
|
||||
/** As above, for the TIMED_OUT_WORKING arm (delivered, but the worker never replies). */
|
||||
@Test
|
||||
void sendTimesOutWorkingNamesFleetPollTargetAndNotTheBareWordRetry() throws Exception {
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null, true));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker starts, never replies
|
||||
|
||||
McpSchema.CallToolResult res = send.get(5, TimeUnit.SECONDS);
|
||||
String text = textOf(res);
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
|
||||
assertTrue(text.contains("fleet_poll{target=\"" + T + "\"}"), "got: " + text);
|
||||
assertFalse(text.contains("retry"), "got: " + text);
|
||||
}
|
||||
|
||||
/** As above, for the BUSY arm (another send already held the session for the whole window). */
|
||||
@Test
|
||||
void sendTimesOutBusyNamesFleetPollTargetAndNotTheBareWordRetry() throws Exception {
|
||||
CompletableFuture<McpSchema.CallToolResult> first = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, T, "first", 5000L, null, Set.of(), null, true));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "the first send should hold the session lock");
|
||||
|
||||
McpSchema.CallToolResult busy = FleetMcp.send(messages, T, "second", 100L, null, Set.of(), null, true);
|
||||
String text = textOf(busy);
|
||||
assertNotEquals(Boolean.TRUE, busy.isError(), "a timeout is informational, not a tool error");
|
||||
assertTrue(text.contains("fleet_poll{target=\"" + T + "\"}"), "got: " + text);
|
||||
assertFalse(text.contains("retry"), "got: " + text);
|
||||
|
||||
// Release the first send so the test thread is not left pinned.
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
FleetMcp.reply(messages, T, Role.WORKER, "first done");
|
||||
first.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #571 (ticket CORRECTION 5): {@code formatReply}'s {@code TIMED_OUT_UNCONFIRMED} arm is
|
||||
* the one message whose whole job is to stop a caller retrying a delivery that may already have
|
||||
@@ -399,7 +459,7 @@ class FleetMcpTest {
|
||||
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
|
||||
herdr.agentSendFailsWith("send_failed");
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null));
|
||||
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null, true));
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
@@ -417,15 +477,59 @@ class FleetMcpTest {
|
||||
+ "a resend here can double-deliver the same brief: got " + text);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #801 defect 1: {@code SEND} is granted to callers {@code DRAIN} refuses (an architect,
|
||||
* a collaborator, an observer), so a TIMED_OUT/BUSY receipt must not name {@code fleet_poll}
|
||||
* when the caller's own grant would be refused there — it must say the reply is unrecoverable
|
||||
* instead.
|
||||
*/
|
||||
@Test
|
||||
void sendTimesOutWithoutDrainGrantNamesNoRouteAndAdmitsItCannotBeRecovered() {
|
||||
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null, false);
|
||||
String text = textOf(res);
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
|
||||
assertFalse(text.contains("fleet_poll{target"), "got: " + text);
|
||||
assertTrue(text.contains("cannot be recovered"), "got: " + text);
|
||||
}
|
||||
|
||||
/** As above, for {@code answer()} — {@code ANSWER} is primary-or-architect, and an architect is refused DRAIN too. */
|
||||
@Test
|
||||
void answerTimesOutWithoutDrainGrantNamesNoRoute() throws Exception {
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null, true));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_a"), "the send must be waiting for the ask to surface to");
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.ask(messages, "term_a", "which config?", 5000L));
|
||||
McpSchema.CallToolResult q = send.get(6, TimeUnit.SECONDS);
|
||||
String qt = textOf(q);
|
||||
String afterMarker = qt.substring(qt.indexOf("turnId=\"") + "turnId=\"".length());
|
||||
String turnId = afterMarker.substring(0, afterMarker.indexOf('"'));
|
||||
|
||||
// The resumed worker never replies, so the answer's own wait times out.
|
||||
McpSchema.CallToolResult answer = FleetMcp.answer(messages, turnId, "config.yaml", 150L, null, false);
|
||||
String text = textOf(answer);
|
||||
assertNotEquals(Boolean.TRUE, answer.isError(), "a timeout is informational, not a tool error");
|
||||
assertFalse(text.contains("fleet_poll{target"), "got: " + text);
|
||||
assertTrue(text.contains("cannot be recovered"), "got: " + text);
|
||||
|
||||
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendRejectsMissingArgs() {
|
||||
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null).isError());
|
||||
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null).isError());
|
||||
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null, true).isError());
|
||||
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null, true).isError());
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
|
||||
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null);
|
||||
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null, true);
|
||||
McpSchema.CallToolResult async = FleetMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
|
||||
|
||||
assertTrue(blocking.isError());
|
||||
@@ -444,7 +548,7 @@ class FleetMcpTest {
|
||||
assertSendRoundTrips("term_live_member", profiles);
|
||||
|
||||
// A herdr-owned pane outside the bridge roster cannot be classified at accept time.
|
||||
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null);
|
||||
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null, true);
|
||||
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
|
||||
}
|
||||
|
||||
@@ -536,7 +640,7 @@ class FleetMcpTest {
|
||||
void askThenAnswerRoundTrips() throws Exception {
|
||||
// The primary delegates and blocks; wait until its waiter is open before the worker asks.
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null));
|
||||
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null, true));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
@@ -558,7 +662,7 @@ class FleetMcpTest {
|
||||
|
||||
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null, true));
|
||||
|
||||
// The worker's ask returns the answer — it resumes the same turn.
|
||||
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
|
||||
@@ -584,7 +688,7 @@ class FleetMcpTest {
|
||||
|
||||
@Test
|
||||
void answerToAStaleTurnIsAnError() {
|
||||
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null);
|
||||
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null, true);
|
||||
assertTrue(res.isError());
|
||||
assertTrue(textOf(res).contains("no longer open"), textOf(res));
|
||||
}
|
||||
|
||||
@@ -266,4 +266,61 @@ class PrimaryRegistryTest {
|
||||
|
||||
assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_never_seen").orElseThrow());
|
||||
}
|
||||
|
||||
// ── fleetd #778: a nudge must offer only what the delegator's own role may run ──────────────
|
||||
|
||||
@Test
|
||||
void theShortDelegationOverloadsGrantEverythingAPrimaryMayRun() {
|
||||
var reg = new PrimaryRegistry(null);
|
||||
reg.recordDelegation("term_worker", "term_lead");
|
||||
|
||||
assertEquals(Boolean.TRUE, reg.nudgeMayDrainFor("term_worker").orElseThrow());
|
||||
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_worker").orElseThrow());
|
||||
}
|
||||
|
||||
@Test
|
||||
void theFullOverloadCarriesTheCallersOwnGrants() {
|
||||
var reg = new PrimaryRegistry(null);
|
||||
reg.recordDelegation("term_worker", "term_observer", "observer-name", false, false);
|
||||
|
||||
assertEquals(Boolean.FALSE, reg.nudgeMayDrainFor("term_worker").orElseThrow());
|
||||
assertEquals(Boolean.FALSE, reg.nudgeMayAnswerFor("term_worker").orElseThrow());
|
||||
}
|
||||
|
||||
@Test
|
||||
void theTwoGrantsAreIndependent() {
|
||||
var reg = new PrimaryRegistry(null);
|
||||
reg.recordDelegation("term_worker", "term_architect", "architect-name", false, true);
|
||||
|
||||
assertEquals(Boolean.FALSE, reg.nudgeMayDrainFor("term_worker").orElseThrow(),
|
||||
"an architect may not DRAIN");
|
||||
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_worker").orElseThrow(),
|
||||
"an architect may ANSWER");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theSingletonPrimaryFallbackIsAlwaysGrantedEverything() {
|
||||
var reg = new PrimaryRegistry("term_pinned");
|
||||
|
||||
assertEquals(Boolean.TRUE, reg.nudgeMayDrainFor("term_never_seen").orElseThrow());
|
||||
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_never_seen").orElseThrow());
|
||||
}
|
||||
|
||||
@Test
|
||||
void withNoDelegationAndNoPrimaryTheGrantsAreUnknown() {
|
||||
var reg = new PrimaryRegistry(null);
|
||||
|
||||
assertTrue(reg.nudgeMayDrainFor("term_never_seen").isEmpty());
|
||||
assertTrue(reg.nudgeMayAnswerFor("term_never_seen").isEmpty());
|
||||
}
|
||||
|
||||
@Test
|
||||
void forgettingADelegationForgetsItsGrantsToo() {
|
||||
var reg = new PrimaryRegistry(null);
|
||||
reg.recordDelegation("term_worker", "term_observer", "observer-name", false, false);
|
||||
reg.forgetDelegation("term_worker");
|
||||
|
||||
assertTrue(reg.nudgeMayDrainFor("term_worker").isEmpty());
|
||||
assertTrue(reg.nudgeMayAnswerFor("term_worker").isEmpty());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -223,4 +223,19 @@ class LeadCoordLoopTest {
|
||||
assertEquals(List.of(), channel.acked());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aDisabledGateDeliversAHeldMessageDespiteUnsubmittedTextInThePromptBox() {
|
||||
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 = new LeadCoordLoop(channel, new AgentControl(herdr), () -> Map.of(LEAD_TERM, SELF),
|
||||
NO_SCHEDULER, 3_000L, false);
|
||||
|
||||
loop.tick();
|
||||
|
||||
assertEquals(1, prompts(herdr).size(), "a disabled gate delivers even though the box holds a draft");
|
||||
assertEquals(List.of("m1"), channel.acked());
|
||||
assertTrue(herdr.calls.stream().noneMatch(c -> c.method().equals("agent.read")),
|
||||
"a disabled gate never reads the pane");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -25,6 +25,7 @@ import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@@ -521,12 +522,20 @@ class LeadHeartbeatLoopTest {
|
||||
private static LeadHeartbeatLoop tickableLoop(FailableHerdrClient herdr, AtomicLong now,
|
||||
List<MemberSession>[] rosterBox, InMemoryReplyInbox inbox,
|
||||
int quietNudgeCap, ScheduledExecutorService scheduler) {
|
||||
return tickableLoop(herdr, now, rosterBox, inbox, quietNudgeCap, scheduler, true);
|
||||
}
|
||||
|
||||
/** As above, with the prompt-box gate's on/off switch. */
|
||||
private static LeadHeartbeatLoop tickableLoop(FailableHerdrClient herdr, AtomicLong now,
|
||||
List<MemberSession>[] rosterBox, InMemoryReplyInbox inbox,
|
||||
int quietNudgeCap, ScheduledExecutorService scheduler,
|
||||
boolean promptBoxGateEnabled) {
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
PrimaryRegistry registry = new PrimaryRegistry(LEAD);
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, agents, inbox, scheduler, 5, 100_000);
|
||||
return new LeadHeartbeatLoop(registry, agents, inbox, () -> rosterBox[0], pushLoop, scheduler,
|
||||
now::get, IDLE_AFTER_NANOS, 100_000L, quietNudgeCap, null,
|
||||
new LeadHeartbeatLoop.LeadContextSource(t -> highReading()), true);
|
||||
new LeadHeartbeatLoop.LeadContextSource(t -> highReading()), true, true, promptBoxGateEnabled);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -695,6 +704,7 @@ class LeadHeartbeatLoopTest {
|
||||
private final String lead;
|
||||
private final List<String> sentTexts = new ArrayList<>();
|
||||
private final List<String> promptTargets = new ArrayList<>();
|
||||
private final AtomicInteger reads = new AtomicInteger();
|
||||
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;
|
||||
@@ -720,6 +730,10 @@ class LeadHeartbeatLoopTest {
|
||||
return List.copyOf(promptTargets);
|
||||
}
|
||||
|
||||
int readCount() {
|
||||
return reads.get();
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
public JsonNode call(String method, Object params) {
|
||||
@@ -730,6 +744,7 @@ class LeadHeartbeatLoopTest {
|
||||
.put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
reads.incrementAndGet();
|
||||
return MAPPER.createObjectNode()
|
||||
.set("read", MAPPER.createObjectNode().put("text", paneTail));
|
||||
}
|
||||
@@ -794,4 +809,23 @@ class LeadHeartbeatLoopTest {
|
||||
"a pane whose box cannot be found may be holding a draft");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aDisabledGateSendsTheNudgeDespiteUnsubmittedTextInThePromptBox() {
|
||||
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, false);
|
||||
|
||||
loop.tick(); // opens the idle window
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
|
||||
loop.tick(); // a disabled gate must not hold on the operator's draft
|
||||
assertEquals(1, herdr.sentTexts().size(),
|
||||
"a disabled gate must send the nudge even though the box holds a draft");
|
||||
assertEquals(0, herdr.readCount(), "a disabled gate must never read the lead's pane");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -107,6 +107,219 @@ class MessageServiceTest {
|
||||
"an unreplied but finished turn resolves via the completion fallback");
|
||||
assertEquals("BUILD GREEN: 391 files", reply.text(), "the scraped transcript tail is returned");
|
||||
assertTrue(reply.completed(), "a scraped completion still counts as completed");
|
||||
assertTrue(messages.drainReplies(T).isEmpty(),
|
||||
"a completion resolved by a live caller must not also land in the inbox");
|
||||
}
|
||||
|
||||
/** Poll {@code messages.drainReplies(target)} until it is non-empty or the deadline passes. */
|
||||
private java.util.List<ReplyInbox.InboxMessage> awaitDrained(String target) throws InterruptedException {
|
||||
var drained = messages.drainReplies(target);
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (drained.isEmpty() && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
drained = messages.drainReplies(target);
|
||||
}
|
||||
return drained;
|
||||
}
|
||||
|
||||
/**
|
||||
* Polls {@code target}'s inbox for {@code millis} and returns {@code true} only if nothing new
|
||||
* ever arrived. Used to assert an absence against an asynchronous completion race: the window
|
||||
* has to be long enough for a virtual thread already scheduled to actually run.
|
||||
*/
|
||||
private boolean noNewReplyWithin(String target, long millis) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + millis;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
if (!messages.drainReplies(target).isEmpty()) {
|
||||
return false;
|
||||
}
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #801: {@code send}'s TimeoutException branch never completes {@code reply} itself, and
|
||||
* {@link Rendezvous#close} only deregisters it from future session lookups — it does not cancel
|
||||
* the future. The CB-106 completion fallback resolves the exact captured future (CB-116), so it
|
||||
* can still succeed after the caller gave up and returned TIMED_OUT_WORKING. With nobody left
|
||||
* awaiting that future, the late scrape must land in the target's inbox instead of being
|
||||
* silently discarded.
|
||||
*/
|
||||
@Test
|
||||
void aCompletionThatArrivesAfterTheCallerGaveUpLandsInTheInbox() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", 150);
|
||||
awaitWaiting();
|
||||
|
||||
herdr.readText("$ prompt"); // pre-turn baseline
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker starts — no reply yet
|
||||
|
||||
MessageService.Reply timedOut = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, timedOut.outcome(),
|
||||
"the caller gives up while the worker is still delivered and working");
|
||||
|
||||
// The worker's turn finishes AFTER the caller already gave up and nobody is listening.
|
||||
herdr.readText("late answer, nobody was listening");
|
||||
injector.onStatus(T, AgentStatus.IDLE); // working -> idle: the completion fallback fires
|
||||
|
||||
var drained = awaitDrained(T);
|
||||
assertEquals(1, drained.size(), "the late completion must land in the inbox, not vanish");
|
||||
assertEquals("late answer, nobody was listening", drained.get(0).content());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #808: once a send has timed out, an explicit {@code fleet_reply} for that same turn
|
||||
* must land in the inbox exactly once — not twice, with the second copy a completion scrape of
|
||||
* the same answer the reply already delivered.
|
||||
*/
|
||||
@Test
|
||||
void anExplicitReplyAfterATimeoutLandsExactlyOnceNotTwice() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", 150);
|
||||
awaitWaiting();
|
||||
|
||||
herdr.readText("$ prompt");
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker starts — no reply yet
|
||||
|
||||
MessageService.Reply timedOut = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, timedOut.outcome(),
|
||||
"the caller gives up while the worker is still delivered and working");
|
||||
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "structured answer"),
|
||||
"no live waiter is open by now, so the explicit reply is held for a later drain");
|
||||
var afterReply = awaitDrained(T);
|
||||
assertEquals(1, afterReply.size(), "the explicit reply must land in the inbox exactly once");
|
||||
assertEquals("structured answer", afterReply.get(0).content());
|
||||
|
||||
// The worker's turn then finishes for real; the completion fallback races the same waiter
|
||||
// the reply already claimed, and must not publish a second entry for it.
|
||||
herdr.readText("late scrape that must not also publish");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertTrue(noNewReplyWithin(T, 500),
|
||||
"a completion scrape for a turn that already got an explicit fleet_reply must not "
|
||||
+ "publish a second entry");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #808: the suppression above is per-turn, not per-target. A second, independent turn on
|
||||
* the same member that never calls {@code fleet_reply} must still have its completion scrape
|
||||
* land, even though an earlier turn on that same target already got an explicit reply.
|
||||
*/
|
||||
@Test
|
||||
void theSuppressionIsPerTurnNotPerTarget() throws Exception {
|
||||
// First turn: times out, then gets an explicit fleet_reply.
|
||||
CompletableFuture<MessageService.Reply> first = sendAsync("first task", 150);
|
||||
awaitWaiting();
|
||||
|
||||
herdr.readText("$ prompt 1");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, first.get(5, TimeUnit.SECONDS).outcome());
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "first answer"));
|
||||
var firstDrain = awaitDrained(T);
|
||||
assertEquals(1, firstDrain.size());
|
||||
assertEquals("first answer", firstDrain.get(0).content());
|
||||
|
||||
// Finishing turn 1 for real also frees the worker for turn 2, and exercises that turn 1's
|
||||
// own late scrape stays suppressed right up to that point.
|
||||
herdr.readText("first turn's late scrape — must stay suppressed");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(noNewReplyWithin(T, 300), "turn 1's own scrape must still be suppressed");
|
||||
|
||||
// Second, independent turn on the same target: times out, and no fleet_reply is ever called.
|
||||
CompletableFuture<MessageService.Reply> second = sendAsync("second task", 150);
|
||||
awaitWaiting();
|
||||
|
||||
herdr.readText("$ prompt 2");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, second.get(5, TimeUnit.SECONDS).outcome());
|
||||
|
||||
herdr.readText("second turn's late scrape — must still land");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
var secondDrain = awaitDrained(T);
|
||||
assertEquals(1, secondDrain.size(),
|
||||
"a later, independent turn's completion must still land even though an earlier "
|
||||
+ "turn on the same target already got an explicit reply");
|
||||
assertEquals("second turn's late scrape — must still land", secondDrain.get(0).content());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #801 defect 2: the publish inside {@code strandLateResolution}'s {@code whenComplete}
|
||||
* callback runs with no caller left to surface a throw to — {@code AmqpReplyInbox.publish} can
|
||||
* throw {@code IllegalStateException}, and an exception thrown inside a {@code whenComplete}
|
||||
* action is captured into the discarded dependent stage, never rethrown. It must be caught and
|
||||
* logged there, or a broker failure on this path is invisible and the late reply is just lost.
|
||||
*/
|
||||
@Test
|
||||
void aLateCompletionThatFailsToPublishIsLoggedNotLost() throws Exception {
|
||||
FakeHerdr localHerdr = new FakeHerdr();
|
||||
Rendezvous localRendezvous = new Rendezvous();
|
||||
java.util.concurrent.atomic.AtomicLong localClock = new java.util.concurrent.atomic.AtomicLong();
|
||||
CompletionResolver localCompletion = new CompletionResolver(new AgentControl(localHerdr), localRendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(),
|
||||
() -> localClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
|
||||
Injector localInjector = new Injector(new AgentControl(localHerdr), localCompletion);
|
||||
ReplyInbox throwingInbox = new ReplyInbox() {
|
||||
private final InMemoryReplyInbox delegate = new InMemoryReplyInbox();
|
||||
|
||||
@Override public void own(String target) { delegate.own(target); }
|
||||
@Override public void release(String target) { delegate.release(target); }
|
||||
|
||||
@Override public void publish(String target, String msgId, String content) {
|
||||
throw new IllegalStateException("PROBE-801-strandLateResolution");
|
||||
}
|
||||
|
||||
@Override public java.util.List<InboxMessage> peek(String target) { return delegate.peek(target); }
|
||||
@Override public boolean ack(String target, String msgId) { return delegate.ack(target, msgId); }
|
||||
};
|
||||
MessageService localMessages = new MessageService(new AgentControl(localHerdr), localInjector,
|
||||
localRendezvous, throwingInbox);
|
||||
|
||||
CompletableFuture<MessageService.Reply> send = CompletableFuture.supplyAsync(
|
||||
() -> localMessages.send(T, "do the task", 150, (String) null));
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!localRendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(localRendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
|
||||
localHerdr.readText("$ prompt");
|
||||
localInjector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
localInjector.onStatus(T, AgentStatus.WORKING); // worker starts — no reply yet
|
||||
|
||||
MessageService.Reply timedOut = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, timedOut.outcome(),
|
||||
"the caller gives up while the worker is still delivered and working");
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
|
||||
try {
|
||||
// The worker's turn finishes AFTER the caller gave up — the completion fallback fires,
|
||||
// and its publish into throwingInbox fails.
|
||||
localHerdr.readText("late answer, publish fails");
|
||||
localInjector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
long logDeadline = System.currentTimeMillis() + 2000;
|
||||
while (appender.list.isEmpty() && System.currentTimeMillis() < logDeadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertFalse(appender.list.isEmpty(), "a failed late-completion publish must be logged");
|
||||
assertEquals(Level.WARN, appender.list.get(0).getLevel());
|
||||
} finally {
|
||||
detachMessageServiceLog(appender);
|
||||
}
|
||||
|
||||
assertFalse(localMessages.hasStrandedReply(T),
|
||||
"strandedReplies must not be set when the publish it depends on failed");
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1034,6 +1247,39 @@ class MessageServiceTest {
|
||||
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link MessageService#ownsTicket(String, String)} is the side-effect-free ownership check
|
||||
* the authorization gate consults ahead of {@link MessageService#poll(String, String)}, which
|
||||
* also fires the ticket's collection hook. Read-only: polling the same ticket afterwards still
|
||||
* sees it, which {@link #theOneArgPollOverloadBypassesOwnershipEntirely} nearby does not
|
||||
* guarantee for every overload.
|
||||
*/
|
||||
@Test
|
||||
void publicOwnsTicketIsASideEffectFreeOwnershipCheck() {
|
||||
Principal lead = Principal.leader("opus", "term_lead", 1);
|
||||
String ticket = messages.sendAsync(T, "long task", null, lead);
|
||||
|
||||
assertTrue(messages.ownsTicket(ticket, lead.ownerKey()));
|
||||
assertFalse(messages.ownsTicket(ticket, Principal.anonymous().ownerKey()));
|
||||
assertFalse(messages.ownsTicket("no-such-ticket", lead.ownerKey()),
|
||||
"a ticket this daemon never heard of is owned by nobody");
|
||||
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket, lead.ownerKey()).phase(),
|
||||
"the ownership check must not have consumed or altered the ticket");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #803: the gap between "unknown ticket" (above) and "no ticket at all". A null ticket
|
||||
* used to reach {@code ConcurrentHashMap.get(null)} and throw a {@code NullPointerException}
|
||||
* instead of returning false.
|
||||
*/
|
||||
@Test
|
||||
void publicOwnsTicketWithANullTicketReturnsFalseRatherThanThrowing() {
|
||||
Principal lead = Principal.leader("opus", "term_lead", 1);
|
||||
assertFalse(messages.ownsTicket(null, lead.ownerKey()),
|
||||
"no ticket at all is owned by nobody, the same as an unknown one");
|
||||
}
|
||||
|
||||
@Test
|
||||
void architectOwnershipUsesTerminalRatherThanSlot() {
|
||||
Principal oldArchitect = Principal.architect("opus", "term_OLD", 1);
|
||||
|
||||
@@ -339,6 +339,85 @@ class ReplyPushLoopTest {
|
||||
assertEquals(0, rec.promptTargets().size(), "nobody live was found, so nothing was ever sent");
|
||||
}
|
||||
|
||||
// --- fleetd #778: never nudge a delegator to run a call its own role is refused -------------
|
||||
|
||||
@Test
|
||||
void onReplyQueuedNeverNudgesADelegatorThatMayNotRunDrain() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
|
||||
|
||||
loop(1, 50).onReplyQueued(WORKER);
|
||||
|
||||
Thread.sleep(200);
|
||||
assertEquals(0, rec.sendCount(),
|
||||
"a delegator whose role cannot run fleet_poll(target=...) must never be nudged to run it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void onReplyQueuedStillNudgesADelegatorThatMayRunDrain() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
registry.recordDelegation(WORKER, OTHER_PRIMARY, "lead", true, true);
|
||||
|
||||
loop(1, 50).onReplyQueued(WORKER);
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"),
|
||||
"a delegator that may run DRAIN must still be nudged about it: " + nudge);
|
||||
}
|
||||
|
||||
@Test
|
||||
void onQuestionOpenedNeverNudgesADelegatorThatMayNotRunAnswer() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
|
||||
|
||||
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||
|
||||
Thread.sleep(200);
|
||||
assertEquals(0, rec.sendCount(),
|
||||
"a delegator whose role cannot run fleet_send(turnId=...) must never be nudged to run it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void onQuestionOpenedStillNudgesADelegatorThatMayRunAnswer() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
registry.recordDelegation(WORKER, OTHER_PRIMARY, "architect", false, true);
|
||||
|
||||
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("fleet_send(turnId="),
|
||||
"an architect may run ANSWER, so the nudge must still name it: " + nudge);
|
||||
}
|
||||
|
||||
/**
|
||||
* A delegator barred from DRAIN can still have a ticket genuinely pending for it — tickets are
|
||||
* out of this unit's scope (fleetd #778 Part 1 already makes a ticket nudge correct for an
|
||||
* observer). The reply half must still be suppressed even though the ticket half fires.
|
||||
*/
|
||||
@Test
|
||||
void aForbiddenReplyNudgeIsSuppressedWhileAnEligibleTicketStillFires() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
|
||||
var loop = loop(1, 300); // wide backoff so both calls land before the first tick fires
|
||||
loop.onReplyQueued(WORKER);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertFalse(nudge.contains("fleet_poll(target=" + WORKER + ")"),
|
||||
"the reply half is forbidden for this delegator and must never be named: " + nudge);
|
||||
}
|
||||
|
||||
// --- nudge format --------------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
@@ -1150,6 +1229,21 @@ class ReplyPushLoopTest {
|
||||
"the nudge lands on the first tick after the box empties");
|
||||
}
|
||||
|
||||
// --- the gate's on/off switch ----------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void aDisabledGateNudgesALeadWithUnsubmittedTextInItsPromptBox() {
|
||||
var herdr = new PaneTextHerdrClient(FakeHerdr.DRAFTED_PROMPT_CARET);
|
||||
agents = new AgentControl(herdr);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(5, 100_000, false);
|
||||
loop.onReplyQueued(WORKER);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
|
||||
"a disabled gate delivers even though the box holds a draft");
|
||||
assertEquals(0, herdr.readCount(), "a disabled gate never reads the pane");
|
||||
}
|
||||
|
||||
// --- helpers -------------------------------------------------------------------------------
|
||||
|
||||
private ReplyPushLoop loop() {
|
||||
@@ -1164,6 +1258,11 @@ class ReplyPushLoopTest {
|
||||
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, metrics);
|
||||
}
|
||||
|
||||
private ReplyPushLoop loop(int maxReminders, long backoffMs, boolean promptBoxGateEnabled) {
|
||||
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, null,
|
||||
promptBoxGateEnabled);
|
||||
}
|
||||
|
||||
private static AgentControl agentWithStatus(String status) {
|
||||
return new AgentControl(new FakeHerdrClient(status));
|
||||
}
|
||||
@@ -1483,6 +1582,7 @@ class ReplyPushLoopTest {
|
||||
private static final class PaneTextHerdrClient implements HerdrClient {
|
||||
private final List<Map.Entry<String, Object>> prompts =
|
||||
Collections.synchronizedList(new ArrayList<>());
|
||||
private final AtomicInteger reads = new AtomicInteger();
|
||||
private volatile String paneText;
|
||||
private volatile boolean failReads = false;
|
||||
volatile CountDownLatch sendLatch = new CountDownLatch(1);
|
||||
@@ -1504,6 +1604,10 @@ class ReplyPushLoopTest {
|
||||
return prompts.size();
|
||||
}
|
||||
|
||||
int readCount() {
|
||||
return reads.get();
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.get".equals(method)) {
|
||||
@@ -1513,6 +1617,7 @@ class ReplyPushLoopTest {
|
||||
.put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.read".equals(method)) {
|
||||
reads.incrementAndGet();
|
||||
if (failReads) {
|
||||
throw new HerdrException("herdr socket read timed out");
|
||||
}
|
||||
|
||||
@@ -285,6 +285,67 @@ class FleetAppAuthTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* An observer may read a ticket its own {@code fleet_send(wait:false)} created, over the REST
|
||||
* route and not only the unit-level classifier — and still not a ticket a different caller
|
||||
* created, even another observer pane's.
|
||||
*/
|
||||
@Test
|
||||
void restPollLetsAnObserverReadItsOwnTicketButNotAnothersOverRest() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, new Rendezvous());
|
||||
|
||||
Javalin observerApp = startOnSharedServiceAsObserver(messages, herdr, 9001L); // -> term_shell
|
||||
Javalin otherObserverApp = startOnSharedServiceAsObserver(messages, herdr, FakeHerdr.WORKER_PID); // -> term_a
|
||||
try {
|
||||
Principal observer = Principal.observer("term_shell", 9001);
|
||||
String ticket = messages.sendAsync("term_a", "long task", null, observer);
|
||||
|
||||
HttpResponse<String> own = send(observerApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, own.statusCode());
|
||||
assertFalse(own.body().contains("forbidden"),
|
||||
"the observer that created the ticket must read it: " + own.body());
|
||||
|
||||
// Unlike a worker's or the primary's mismatched ticket (refused by MessageService's own
|
||||
// ownership check, inside a 200 response), a different observer is refused at the
|
||||
// authorization gate itself, since its grant is conditional on the classifier -- so
|
||||
// this is a 403, never reaching poll().
|
||||
HttpResponse<String> other = send(otherObserverApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(403, other.statusCode(), other.body());
|
||||
assertFalse(other.body().contains("\"reply\""),
|
||||
"a refusal must never carry reply text: " + other.body());
|
||||
} finally {
|
||||
observerApp.stop();
|
||||
otherObserverApp.stop();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #startOnSharedService(MessageService, FakeHerdr, long)}, but resolves every
|
||||
* connecting pane to the unconfigured-pane {@link Role#OBSERVER} floor instead of a spawned
|
||||
* worker — no lead, collaborator, or member map claims it.
|
||||
*/
|
||||
private Javalin startOnSharedServiceAsObserver(MessageService messages, FakeHerdr herdr, long pid) {
|
||||
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
|
||||
agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
Map.of(wcfg.profile(), wcfg), wcfg.profile(),
|
||||
k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
|
||||
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, new MemberRegistry(null), t -> null, Map::of);
|
||||
Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
|
||||
|
||||
return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
|
||||
callers, appMetrics).build().start("127.0.0.1", 0);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
|
||||
* caller's own terminal on the ticket it returns, so that caller can still poll its own
|
||||
@@ -745,6 +806,59 @@ class FleetAppAuthTest {
|
||||
return app.port();
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #startWithRealClassifier}, but with a {@link MemberRegistry} bound to an architect
|
||||
* slot at {@code terminal} (the defense-in-depth path, no live roster) — lets a test drive a
|
||||
* real architect {@link Principal} over HTTP.
|
||||
*/
|
||||
private int startWithArchitect(long pid, String terminal) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
|
||||
agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
Map.of(wcfg.profile(), wcfg), wcfg.profile(),
|
||||
k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
|
||||
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, new Rendezvous());
|
||||
|
||||
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", terminal));
|
||||
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, Map::of, members);
|
||||
metrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
|
||||
|
||||
app = new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
|
||||
callers, metrics).build().start("127.0.0.1", 0);
|
||||
return app.port();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #801 defect 1: {@code SEND} is granted to an architect but {@code DRAIN} is not, so a
|
||||
* TIMED_OUT/BUSY receipt to an architect must not name {@code GET /sessions/{id}/replies} —
|
||||
* that route refuses the same caller — and must say the reply is unrecoverable instead.
|
||||
*/
|
||||
@Test
|
||||
void anArchitectsTimedOutSendDoesNotNameTheRepliesRouteItCannotDrain() throws Exception {
|
||||
int port = startWithArchitect(FakeHerdr.WORKER_PID, "term_a"); // FakeHerdr resolves this pid to term_a
|
||||
|
||||
HttpResponse<String> res = send(port, "POST", "/sessions/term_b/message",
|
||||
"{\"content\":\"hi\",\"timeoutMs\":150}", null);
|
||||
assertEquals(202, res.statusCode(), "an architect is granted SEND: " + res.body());
|
||||
JsonNode body = new ObjectMapper().readTree(res.body());
|
||||
String detail = body.get("detail").asText();
|
||||
assertFalse(detail.contains("/replies"), "got: " + detail);
|
||||
assertTrue(detail.contains("cannot be recovered"), "got: " + detail);
|
||||
|
||||
assertEquals(403, send(port, "GET", "/sessions/term_b/replies", null, null).statusCode(),
|
||||
"the architect this receipt was sent to is indeed refused DRAIN");
|
||||
}
|
||||
|
||||
// --- loopback-trust: the caller is the primary -------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -663,8 +663,10 @@ class FleetAppTest {
|
||||
|
||||
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":150}");
|
||||
assertEquals(202, res.statusCode());
|
||||
assertEquals("queued", mapper.readTree(res.body()).get("status").asText());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("queued", body.get("status").asText());
|
||||
assertFalse(herdr.called("agent.prompt"), "no injection while the worker is mid-turn");
|
||||
assertDetailNamesTheInboxRouteNotARetry(body);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -674,8 +676,21 @@ class FleetAppTest {
|
||||
|
||||
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":250}");
|
||||
assertEquals(202, res.statusCode());
|
||||
assertEquals("working", mapper.readTree(res.body()).get("status").asText());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("working", body.get("status").asText());
|
||||
assertTrue(herdr.called("agent.prompt"), "message was injected");
|
||||
assertDetailNamesTheInboxRouteNotARetry(body);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #801: the TIMED_OUT_WORKING / TIMED_OUT_QUEUED / BUSY receipt is never tracked by a
|
||||
* ticket, so it must name the one route that recovers a late reply — draining the session's
|
||||
* inbox — and must not invite a bare resend, which can duplicate the delivery.
|
||||
*/
|
||||
private void assertDetailNamesTheInboxRouteNotARetry(JsonNode body) {
|
||||
String detail = body.get("detail").asText();
|
||||
assertTrue(detail.contains("/sessions/term_a/replies"), "got: " + detail);
|
||||
assertFalse(detail.contains("retry"), "got: " + detail);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user