Compare commits

..

7 Commits

Author SHA1 Message Date
Dai Ha a3eeace447 handover skill: BOOTSTRAP_NEVER_SENT, and relaunchReadySeconds bounds three waits
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 1m2s
CI / build (push) Failing after 1m55s
Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-06 19:28:21 +02:00
Dai Ha b696c31756 fleetd #796: retry bootstrap delivery to a fresh lead pane
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m5s
CI / build (push) Failing after 2m17s
A roll killed the old pane, launched the new one, then lost bootstrapText to a
herdr agent_not_ready refusal. The successor woke with no handover and no way to
tell a failed roll from a cold start.

sendBootstrapWithRetry retries only that refusal, bounded by
relaunchReadySeconds. A persistent refusal now ends the roll with the new
BOOTSTRAP_NEVER_SENT outcome, so fleet_handover{action:"status"} can report it.

Verified at 027d413 in a throwaway worktree: 2210 tests, 0 failures, 0 errors,
clean install.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-06 19:27:36 +02:00
Dai Ha 96ebae28a6 reviewer skill: send fleet_reply before writing the reasoning out
Five reviewer turns were lost in one session. Each wrote a long, correct
analysis to its own pane and ended the turn with no fleet_reply, so the
bridge scraped the pane and the lead received a clipped fragment.

The briefs are half the cause: each carried a five-item checklist of things
to hunt alongside a ~90-word capped output format, which reads as two
contradictory output contracts. The skill now says the checklist is where to
look, not the shape of the answer, and that the reply goes out as soon as
the answer is known.

Measured: tickets task-7da785-29 and -30 resolved at 18:28 via the
turn-completion fallback, 1233 and 1403 chars scraped.
2026-10-06 19:01:38 +02:00
Dai Ha b41aa663f7 Bridge block: invariant 6 — an operator outranks the confirm rule
Invariant 6 as first written told every receiver to confirm, with no
mention of who may override it. A peer cannot impose a rule on another
operator's session.

Measured: the trinotes pane (observer, work config dir) answered the
communication test and then said its operator's standing rule is not to
answer fleet messages, that a peer cannot change that rule, and that it
would ask its operator before following mine. That is the correct reading
and the rule now says so.

Also records the symmetric error: a sender must not read silence as
agreement or as a dead session.

wiki/7-Use-Cases.md synced byte-identical.
2026-10-06 18:48:12 +02:00
Dai Ha 027d413ce9 fleetd #796: document bootstrap retry readiness bound
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m59s
2026-10-06 18:26:01 +02:00
Dai Ha f544906621 Bridge block: a receiver confirms what it receives
Adds invariant 6 to the canonical block. Nothing in it said a received
message must be answered, so a sender could not tell a handled message
from one that never arrived.

Measured today: the anki pane (observer) sent to vms, which reads
deliverable:false because it has not contacted the daemon since the 08:05
restart. The send was accepted, held 83s, and failed with the body 'nst' --
three characters scraped off the vms screen. The sender read that as an
inconclusive result and asked the lead what happened. fleetd #757 covers
the accept-time refusal; this rule covers the half the protocol owns.

Operator's words: 'when a message sent, at least the receiver should
confirm unless explicit told to not reply'.

wiki/7-Use-Cases.md template synced byte-identical (wiki 3ad606f).
2026-10-06 18:18:10 +02:00
Dai Ha c88f01ecb8 fleetd #796: retry bootstrap delivery after agent readiness refusal
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 2m4s
2026-10-06 18:16:47 +02:00
19 changed files with 214 additions and 544 deletions
+4 -2
View File
@@ -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
+9
View File
@@ -37,6 +37,15 @@ confident-but-wrong finding. Anything you could settle by reading more code is y
## 4. The finding — what goes in `fleet_reply`
**Call `fleet_reply` as soon as you know your answer, before you write the reasoning out.** The
four lines below are the whole deliverable, and they are short on purpose. Analysis you type into
your terminal reaches nobody: when a turn ends with no `fleet_reply`, the bridge scrapes the pane
and the lead receives a clipped fragment instead of a finding. A long, correct analysis and no
`fleet_reply` is a failed turn, and it is the most common way this role fails.
If the delegation also handed you a list of things to check, that list is where to *look*. It is
not the shape of the answer. Work the list, then still send these four lines.
Report the **single most important** real issue in the scope, in these four lines, under
~90 words:
+12
View File
@@ -95,6 +95,18 @@ and the sender silently receives nothing. Fail toward the recoverable error.
5. **Never move a fleet session, pane or peer except through the bridge.** The bridge owns policy;
the multiplexer owns PTYs. Any route that changes fleet state without the bridge's checks
bypasses every rule above — the `herdr` CLI and its socket are the usual example.
6. **Confirm what you receive.** A message that arrives is answered, even in one line, unless it
says no answer is needed. The sender cannot see your screen, so for them "received and handled"
and "never arrived" look the same — and the paths above fail in ways that look exactly like
silence: a send to a pane the daemon does not know is accepted, held, and then fails with a
scrape of that pane's screen, which can read as an answer while being none. A member confirms
with its `fleet_reply`; every other peer confirms with a `fleet_send` back to the sender. If you
cannot do the thing asked, say that — a refusal is a confirmation. **Never read a failed
ticket's body as a reply.** Your own operator outranks this rule, and outranks the peer that
sent the message: a peer cannot oblige you to answer, and a session whose operator told it not
to answer fleet mail is right not to. Where you can, say that much and nothing more. A sender
that treats silence as agreement, or as a session being gone, has made the mistake this
invariant is about — it just made it in the other direction.
### Primary (lead) — run this on every task, in order
@@ -86,52 +86,29 @@ public final class Authz {
*/
public static final Predicate<String> NO_OBSERVER_SEND_TARGET = target -> false;
/**
* The fail-closed classifier for an observer's {@code TASK_READ}: answers no for every
* ticket, so the grant is refused unless a caller supplies a real one. {@code
* MessageService#ownsTicket(String, String)} is the real one, so a ticket this classifier
* accepts is one that caller's own {@code fleet_send(wait:false)} actually created.
*/
public static final Predicate<String> NO_OBSERVER_OWNED_TICKET = ticket -> false;
/**
* Convenience form for a caller with no classifier to supply. Fails closed: a collaborator's
* or an observer's {@code SEND}, and an observer's {@code TASK_READ}, are refused, as if no
* terminal were a configured lead, collaborator, or observer-reachable target, and no ticket
* were the caller's own — the same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR},
* {@link #NO_OBSERVER_SEND_TARGET} and {@link #NO_OBSERVER_OWNED_TICKET} give explicitly.
* Every other action's result is identical to the full form's, since none of them consult
* any of the three classifiers.
* 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.
*
* <p>Its default classifiers deny every collaborator and every observer, so a caller
* enforcing authorization must use the full form instead.
* enforcing authorization must use the five-argument form instead.
*/
public static boolean permits(Principal caller, Action action, String targetSession) {
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR,
NO_OBSERVER_SEND_TARGET, NO_OBSERVER_OWNED_TICKET);
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR, NO_OBSERVER_SEND_TARGET);
}
/**
* As {@link #permits(Principal, Action, String)}, with a real classifier for a collaborator's
* {@code SEND}. An observer's {@code SEND} and {@code TASK_READ} still fail closed — a
* caller enforcing all three grants must use the full form.
* {@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.
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator) {
return permits(caller, action, targetSession, knownLeadOrCollaborator,
NO_OBSERVER_SEND_TARGET, NO_OBSERVER_OWNED_TICKET);
}
/**
* As {@link #permits(Principal, Action, String, Predicate)}, with a real classifier for an
* observer's {@code SEND} too. An observer's {@code TASK_READ} still fails closed — a caller
* enforcing all three grants must use the full form.
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> observerSendTarget) {
return permits(caller, action, targetSession, knownLeadOrCollaborator, observerSendTarget,
NO_OBSERVER_OWNED_TICKET);
return permits(caller, action, targetSession, knownLeadOrCollaborator, NO_OBSERVER_SEND_TARGET);
}
/**
@@ -139,10 +116,9 @@ 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}, for an
* observer's {@code SEND}, and — carrying a ticket id instead
* of a session id — for an observer's {@code TASK_READ},
* ignored otherwise, may be {@code null}
* {@code INBOX}), for a collaborator's {@code SEND}, and for an
* observer's {@code SEND}, ignored otherwise, may be
* {@code null}
* @param knownLeadOrCollaborator whether a terminal is a configured lead or collaborator —
* consulted only for a collaborator's {@code SEND}, to confine
* it to another named peer and never a spawned member's
@@ -152,16 +128,10 @@ 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 observerOwnsTicket whether {@code targetSession} (here, a ticket id) was created
* by this same caller — consulted only for an observer's
* {@code TASK_READ}, to confine it to a ticket its own
* {@code fleet_send(wait:false)} created, never another
* caller's
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> observerSendTarget,
Predicate<String> observerOwnsTicket) {
Predicate<String> observerSendTarget) {
if (caller == null || caller.isAnonymous()) {
return false; // authenticated as nothing ⇒ authorized for nothing
}
@@ -210,15 +180,11 @@ 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 — with one exception: an observer may poll a ticket its
// own fleet_send(wait:false) created, confined by observerOwnsTicket. Session status
// is not ticket-scoped, so that call site supplies no real classifier here and an
// observer's TASK_READ on it stays refused. MessageService compares a ticket's
// creator to the caller on every read as well, so dropping this gate would not expose
// another session's reply — it would move the refusal later and widen what a caller
// that never orchestrates can probe.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect()
|| (caller.isObserver() && observerOwnsTicket.test(targetSession));
// 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();
// 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
@@ -51,11 +51,10 @@ 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, {@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}).
* {@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}).
*/
OBSERVER,
@@ -1466,7 +1466,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 +1479,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
@@ -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;
@@ -528,13 +528,8 @@ 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;
// 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);
Runnable onAccepted = () ->
primaryRegistry.recordDelegation(target, callerTerminal, delegatorName);
// 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)
@@ -579,18 +574,10 @@ 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.
Authz.Action action = toolAction("fleet_poll", a);
// TASK_READ reads a ticket, not a DRAIN target -- the gate must see the
// ticket id, and an observer's grant is confined to a ticket it created.
String authzTarget = action == Authz.Action.TASK_READ ? ticket : target;
Predicate<String> observerOwnsTicket = action == Authz.Action.TASK_READ
? t -> messages.ownsTicket(t, principal(exchange).ownerKey())
: Authz.NO_OBSERVER_OWNED_TICKET;
McpSchema.CallToolResult denied = deny(exchange, action, authzTarget, observerOwnsTicket);
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
if (denied != null) return denied;
return poll(messages, leadChannel, ticket, target, coordId,
return poll(messages, leadChannel, str(a, "ticket"), target, coordId,
principal(exchange).ownerKey());
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
@@ -758,16 +745,6 @@ public final class FleetMcp {
return denyFor(principal(exchange), action, target);
}
/**
* As {@link #deny(McpSyncServerExchange, Authz.Action, String)}, with a real classifier for
* an observer's {@code TASK_READ} — the only call site that can supply one is a ticket poll,
* which knows the ticket id and the caller's owner key.
*/
private McpSchema.CallToolResult deny(McpSyncServerExchange exchange, Authz.Action action,
String target, Predicate<String> observerOwnsTicket) {
return denyFor(principal(exchange), action, target, observerOwnsTicket);
}
/**
* The policy half of {@link #deny}: everything except pulling the caller out of the MCP
* exchange. Kept separate so the authorization decision — the actual control — is unit-testable
@@ -781,22 +758,13 @@ public final class FleetMcp {
* @return {@code null} when the call may proceed, or the error result to return when it may not
*/
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target) {
return denyFor(caller, action, target, Authz.NO_OBSERVER_OWNED_TICKET);
}
/**
* As {@link #denyFor(Principal, Authz.Action, String)}, with a real classifier for an
* observer's {@code TASK_READ}.
*/
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target,
Predicate<String> observerOwnsTicket) {
// The enforcement switch lives HERE rather than in the exchange-facing wrapper: any future
// tool that calls this directly must not be able to skip the gate by accident.
if (!authorizationEnforced) {
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
}
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator(),
callers.observerSendTarget(), observerOwnsTicket)) {
callers.observerSendTarget())) {
if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
}
@@ -48,13 +48,8 @@ public final class PrimaryRegistry {
*/
private final ConcurrentHashMap<String, Delegation> leadByTarget = new ConcurrentHashMap<>();
/**
* 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) {
/** A recorded delegator: the terminal learned from call traffic, and its name, if it has one. */
private record Delegation(String terminal, String name) {
}
/**
@@ -133,28 +128,14 @@ public final class PrimaryRegistry {
/**
* As {@link #recordDelegation(String, String)}, additionally recording the delegating lead's
* 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.
* name when the caller carries one. See {@link #record(String, String)} for why the name
* matters.
*/
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), mayDrainNudge, mayAnswerNudge));
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName)));
}
/** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */
@@ -186,32 +167,6 @@ 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
@@ -1492,16 +1492,6 @@ public final class MessageService {
|| Objects.equals(callerOwner, task.creatorOwner);
}
/**
* Whether {@code callerOwner} owns {@code ticket}, with no side effect — unlike
* {@link #poll(String, String)}, this never fires the ticket's collection hook. A ticket this
* daemon has never heard of is owned by nobody.
*/
public boolean ownsTicket(String ticket, String callerOwner) {
Task task = tasks.get(ticket);
return task != null && ownsTicket(task, callerOwner);
}
/**
* Test seam only — carries no production behaviour, and nothing in this class calls it;
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
@@ -178,20 +178,13 @@ 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 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) {
/** A pending reply target: which lead to nudge, and how many nudges have named it so far. */
private record ReplyEntry(String lead, int nudgeCount) {
}
/**
@@ -214,12 +207,11 @@ 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, 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.
* 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).
*/
private record PendingQuestion(String turnId, String ticket, String target, String lead,
String question, int nudgeCount, boolean mayAnswerNudge) {
String question, int nudgeCount) {
}
private record IncidentLead(String incidentId, String lead) {
@@ -237,11 +229,7 @@ public final class ReplyPushLoop {
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
private List<PendingQuestion> pendingQuestionsFor(String lead) {
// 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();
return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList();
}
/** Question turnIds still open for {@code lead} — a plain snapshot for race comparison. */
@@ -397,7 +385,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 nudge target {}, stopping reminder", lead);
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
return Action.STOP;
}
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
@@ -406,7 +394,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 nudge target {} on every source with pending work, stopping",
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
@@ -415,15 +403,15 @@ public final class ReplyPushLoop {
try {
status = agents.status(lead);
} catch (RuntimeException e) {
log.debug("push: status check failed for nudge target {}, will retry", lead, e);
log.debug("push: status check failed for lead {}, will retry", lead, e);
return Action.WAIT_BUSY;
}
if (!status.injectable()) {
log.debug("push: nudge target {} is {} (not injectable), waiting", lead, status);
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
return Action.WAIT_BUSY;
}
if (!promptBox.clearToSubmit(lead)) {
log.debug("push: nudge target {} has unsubmitted text in its prompt box, waiting", lead);
log.debug("push: lead {} has unsubmitted text in its prompt box, waiting", lead);
return Action.WAIT_BUSY;
}
return Action.INJECT;
@@ -478,14 +466,14 @@ public final class ReplyPushLoop {
if (lead.isEmpty() || isLive(lead.get())) {
return lead;
}
log.debug("push: nudge target {} delegated to for {} is no longer live, forgetting the stale binding "
log.debug("push: lead {} 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 nudge target {} for {} is also not live, skipping this tick", fallback.get(), target);
log.debug("push: fallback lead {} for {} is also not live, skipping this tick", fallback.get(), target);
return Optional.empty();
}
@@ -508,9 +496,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: nudge target {} no longer exists ({})", lead, e.toString());
log.debug("push: lead {} no longer exists ({})", lead, e.toString());
} else {
log.debug("push: liveness check for nudge target {} was inconclusive ({}); treating as live "
log.debug("push: liveness check for lead {} was inconclusive ({}); treating as live "
+ "rather than risk destroying a live binding", lead, e.toString());
}
return !gone;
@@ -529,12 +517,11 @@ public final class ReplyPushLoop {
public void onReplyQueued(String target) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.debug("push: no nudge target is known to be waiting on {}, skipping reminder", target);
log.debug("push: no lead 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(), mayDrainNudge));
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
startOrCoalesce(lead.get());
}
@@ -558,7 +545,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 nudge target is known to be waiting on ticket {} (target {}), skipping nudge",
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
ticket, target);
return;
}
@@ -593,13 +580,11 @@ 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 nudge target is known to be waiting on {}'s question (turnId {}), skipping nudge",
log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge",
target, turnId);
return;
}
boolean mayAnswerNudge = primaryRegistry.nudgeMayAnswerFor(target).orElse(false);
pendingQuestions.put(turnId,
new PendingQuestion(turnId, ticket, target, lead.get(), question, 0, mayAnswerNudge));
pendingQuestions.put(turnId, new PendingQuestion(turnId, ticket, target, lead.get(), question, 0));
startOrCoalesce(lead.get());
}
@@ -623,7 +608,7 @@ public final class ReplyPushLoop {
for (String target : targets) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.warn("push: backend incident {} has no known nudge target for worker {}", incidentId, target);
log.warn("push: backend incident {} has no known lead for target {}", incidentId, target);
continue;
}
targetsByLead.computeIfAbsent(lead.get(), _ -> new ArrayList<>()).add(target);
@@ -659,10 +644,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 nudge target {}, work coalesced in", lead);
log.debug("push: reminder loop already active for lead {}, work coalesced in", lead);
return;
}
log.debug("push: starting reminder loop for nudge target {}", lead);
log.debug("push: starting reminder loop for lead {}", lead);
scheduleNext(lead);
}
@@ -767,11 +752,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 nudge target {} raced the reminder loop's stop — restarting", lead);
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
scheduleNext(lead);
return;
}
log.debug("push: reminder loop ended for nudge target {}", lead);
log.debug("push: reminder loop ended for lead {}", lead);
}
/** Send one combined nudge covering everything currently pending for {@code lead}. */
@@ -788,13 +773,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 nudge target {} drained before the nudge could be sent", lead);
log.debug("push: pending work for lead {} 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 nudge target {} (reply {}/{}, ticket {}/{}, question {}/{}; "
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
+ "{} reply target(s), {} ticket(s), {} question(s))",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders,
@@ -811,7 +796,7 @@ public final class ReplyPushLoop {
}
}
} catch (RuntimeException e) {
log.warn("push: failed to nudge target {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders, e.toString());
}
@@ -829,8 +814,7 @@ 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, e.mayDrainNudge()));
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
}
for (PendingTicket ticket : tickets) {
pendingTickets.computeIfPresent(ticket.ticket(),
@@ -839,7 +823,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.mayAnswerNudge()));
e.nudgeCount() + 1));
}
for (PendingIncident incident : incidents) {
pendingIncidents.computeIfPresent(incident.key(), (id, e) -> new PendingIncident(e.key(),
@@ -293,20 +293,6 @@ 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 {@code TASK_READ} is checked against; pass a real
* {@code MessageService#ownsTicket(String, String)}-backed predicate to exercise the real
* production gate for a ticket poll, as {@link #allow} does.
*/
static boolean permitsFor(Principal caller, Authz.Action action, String target,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> observerSendTarget,
Predicate<String> observerOwnsTicket) {
return Authz.permits(caller, action, target, knownLeadOrCollaborator, observerSendTarget,
observerOwnsTicket);
}
/**
* 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}.
@@ -316,22 +302,11 @@ public final class FleetApp {
* yours" (a worker reaching for another worker's session, or for orchestration).
*/
private boolean allow(Context ctx, Authz.Action action, String target) {
return allow(ctx, action, target, Authz.NO_OBSERVER_OWNED_TICKET);
}
/**
* As {@link #allow(Context, Authz.Action, String)}, with a real classifier for an observer's
* {@code TASK_READ} — the only call site that can supply one is a ticket poll, which knows the
* ticket id and the caller's owner key.
*/
private boolean allow(Context ctx, Authz.Action action, String target,
Predicate<String> observerOwnsTicket) {
if (auth == null) {
return true; // legacy: authorization not enforced
}
Principal caller = ctx.attribute(CALLER);
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.observerSendTarget(),
observerOwnsTicket)) {
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.observerSendTarget())) {
if (action != Authz.Action.READ && action != Authz.Action.METRICS
&& action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
@@ -953,15 +928,11 @@ public final class FleetApp {
/** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */
private void taskStatus(Context ctx) {
String ticket = ctx.pathParam("ticket");
Principal caller = ctx.attribute(CALLER);
String callerOwner = caller == null ? null : caller.ownerKey();
// The gate must see the ticket id, not a literal null -- an observer's grant is confined
// to a ticket its own caller created.
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), ticket, t -> messages.ownsTicket(t, callerOwner))) {
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
return;
}
MessageService.TaskView v = messages.poll(ticket, callerOwner);
Principal caller = ctx.attribute(CALLER);
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.ownerKey());
if (v == null) {
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
return;
@@ -310,17 +310,17 @@ class AuthzTest {
}
/**
* 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.
* 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.
*/
@Test
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxSendAndTaskRead() {
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxAndSend() {
for (Authz.Action a : Authz.Action.values()) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND
|| a == TASK_READ) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND) {
continue;
}
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
@@ -328,49 +328,6 @@ class AuthzTest {
}
}
/**
* {@code TASK_READ} is conditional for an observer, not fixed: it may poll a ticket its own
* {@code fleet_send(wait:false)} created, confined by {@code observerOwnsTicket}, and nothing
* else — a session's status is not ticket-scoped, so no call site can ever supply a predicate
* that grants it, and the 3-/4-/5-argument convenience forms stay fail-closed for it too.
*/
@Test
void anObserverMayTaskReadOnlyATicketTheClassifierSaysItOwns() {
assertTrue(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
"the classifier accepting the ticket as the caller's own must grant TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "someone_elses_ticket",
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, Authz.NO_OBSERVER_SEND_TARGET, target -> false),
"the classifier refusing the ticket must deny TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket"),
"the three-argument convenience form fails closed, so TASK_READ is refused without "
+ "a real classifier");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR),
"the four-argument form still fails closed for TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET),
"the five-argument form still fails closed for TASK_READ, since it supplies no "
+ "observerOwnsTicket classifier either");
}
/**
* Control for the test above: every other action's result for an observer does not move when
* the ticket-ownership classifier does. Only {@code TASK_READ} is wired to it.
*/
@Test
void theTicketOwnershipClassifierMovesOnlyTaskReadForAnObserver() {
for (Authz.Action a : Authz.Action.values()) {
if (a == TASK_READ) {
continue;
}
assertEquals(
Authz.permits(OBSERVER, a, "term_observer"),
Authz.permits(OBSERVER, a, "term_observer", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
a + " must not depend on the ticket-ownership classifier at all");
}
}
@Test
void anObserverIsNotCountedAsAnyOtherRole() {
assertFalse(OBSERVER.isPrimary());
@@ -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,7 +56,6 @@ class FleetMcpAuthzTest {
private final AgentControl agents = new AgentControl(herdr);
private Metrics metrics;
private FleetMcp mcp;
private MessageService messages;
@AfterEach
void close() {
@@ -83,7 +82,7 @@ class FleetMcpAuthzTest {
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> "tok");
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
messages = new MessageService(agents, new Injector(agents), new Rendezvous(),
MessageService 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());
@@ -347,37 +346,6 @@ 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");
}
/**
* {@code fleet_status} stays refused for an observer even once {@code TASK_READ} opens: it is
* not ticket-scoped, so no call site supplies a real {@code observerOwnsTicket} classifier,
* and the default 3-argument {@code denyFor} fails closed.
*/
@Test
void anObserverMayNotReadSessionStatusEvenAfterTaskReadOpensForTickets() {
FleetMcp m = mcp(true);
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, "term_a"),
"a session id is not a ticket the caller could ever own");
}
@Test
void theLegacyConstructorLeavesTheGateOpen() {
// The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced.
@@ -266,61 +266,4 @@ 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());
}
}
@@ -1034,27 +1034,6 @@ class MessageServiceTest {
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
}
/**
* {@link MessageService#ownsTicket(String, String)} is the side-effect-free ownership check
* the authorization gate consults ahead of {@link MessageService#poll(String, String)}, which
* also fires the ticket's collection hook. Read-only: polling the same ticket afterwards still
* sees it, which {@link #theOneArgPollOverloadBypassesOwnershipEntirely} nearby does not
* guarantee for every overload.
*/
@Test
void publicOwnsTicketIsASideEffectFreeOwnershipCheck() {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "long task", null, lead);
assertTrue(messages.ownsTicket(ticket, lead.ownerKey()));
assertFalse(messages.ownsTicket(ticket, Principal.anonymous().ownerKey()));
assertFalse(messages.ownsTicket("no-such-ticket", lead.ownerKey()),
"a ticket this daemon never heard of is owned by nobody");
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket, lead.ownerKey()).phase(),
"the ownership check must not have consumed or altered the ticket");
}
@Test
void architectOwnershipUsesTerminalRatherThanSlot() {
Principal oldArchitect = Principal.architect("opus", "term_OLD", 1);
@@ -339,85 +339,6 @@ class ReplyPushLoopTest {
assertEquals(0, rec.promptTargets().size(), "nobody live was found, so nothing was ever sent");
}
// --- fleetd #778: never nudge a delegator to run a call its own role is refused -------------
@Test
void onReplyQueuedNeverNudgesADelegatorThatMayNotRunDrain() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
loop(1, 50).onReplyQueued(WORKER);
Thread.sleep(200);
assertEquals(0, rec.sendCount(),
"a delegator whose role cannot run fleet_poll(target=...) must never be nudged to run it");
}
@Test
void onReplyQueuedStillNudgesADelegatorThatMayRunDrain() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "lead", true, true);
loop(1, 50).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"),
"a delegator that may run DRAIN must still be nudged about it: " + nudge);
}
@Test
void onQuestionOpenedNeverNudgesADelegatorThatMayNotRunAnswer() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
Thread.sleep(200);
assertEquals(0, rec.sendCount(),
"a delegator whose role cannot run fleet_send(turnId=...) must never be nudged to run it");
}
@Test
void onQuestionOpenedStillNudgesADelegatorThatMayRunAnswer() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
registry.recordDelegation(WORKER, OTHER_PRIMARY, "architect", false, true);
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("fleet_send(turnId="),
"an architect may run ANSWER, so the nudge must still name it: " + nudge);
}
/**
* A delegator barred from DRAIN can still have a ticket genuinely pending for it — tickets are
* out of this unit's scope (fleetd #778 Part 1 already makes a ticket nudge correct for an
* observer). The reply half must still be suppressed even though the ticket half fires.
*/
@Test
void aForbiddenReplyNudgeIsSuppressedWhileAnEligibleTicketStillFires() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
var loop = loop(1, 300); // wide backoff so both calls land before the first tick fires
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertFalse(nudge.contains("fleet_poll(target=" + WORKER + ")"),
"the reply half is forbidden for this delegator and must never be named: " + nudge);
}
// --- nudge format --------------------------------------------------------------------------
@Test
@@ -285,67 +285,6 @@ 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