Compare commits

...

6 Commits

Author SHA1 Message Date
Dai Ha f4176ae455 fleetd #737 unit 3: probe the fallback terminal and resolve it by lead name
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m55s
ReplyPushLoop.resolveLiveLead dropped a dead per-target delegation and
retried PrimaryRegistry.nudgeTargetFor, but returned that fallback without
checking isLive. The fallback is now probed the same way the first lead is,
and the method returns empty rather than trust a dead terminal.

PrimaryRegistry records a delegating lead's name alongside its learned
terminal and resolves the name back to its current terminal at nudge time,
through a name-to-terminal lookup backed by the live lead-tab scan. A name
with no current match falls back to the terminal that was actually learned,
so an unnamed primary, an off-host lead, or a non-herdr lead keeps working
exactly as before. LeadHeartbeatLoop now reads the resolved current terminal
instead of the raw learned one.
2026-10-05 05:29:06 +02:00
Dai Ha 8fa8d18c97 fleetd #726: rewrite the handover skill for the restart roll
The roll now ends the pane and launches a fresh process instead of typing
/clear, and that jar is deployed, so the skill's self-dated "a change is
coming" note had to go.

Measured on 2026-10-05 before writing:

  grep -c "lead-rollover: rolled" fleetd/fleetd.out   -> 20
  grep -c "lead-rollover:" fleetd/fleetd.out          -> 86 (control)
  <tail from the last "fleetd listening"> | grep -c "lead-rollover:" -> 0
  grep -c 'RELAUNCH_NEVER_READY\|RELAUNCH_NOT_RECOGNISED\|OLD_PANE_NEVER_DIED' -> 0

So all 20 recorded rolls ran under /clear and the restart path has never
executed. The section says that rather than implying the old numbers
describe it.

Also:
- name all eight RollState outcomes, with what each one guarantees
- state that relaunchReadySeconds bounds each of two waits, not the pair
- drop the "never observed as WORKING after 8 consecutive IDLE/DONE polls"
  paragraph: grep finds that wait is deleted, so it cannot appear
- drop the #621 warning: contextNotice now takes requireOperatorConfirm
- split the surprise bullets into their own section
2026-10-05 05:04:52 +02:00
Dai Ha 9d653e86df fleetd #726: name RELAUNCH_NEVER_READY in the IN_PROGRESS terminal-state list
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 2m7s
The javadoc on RollState.IN_PROGRESS enumerates the terminal states the entry
can be overwritten with, and omitted RELAUNCH_NEVER_READY. That state is
reachable at LeadRollover.java:710, so the list told a reader a state could not
occur when it can. Found by a reviewer on PR #742, outside its assigned scope.

Comment only; no behaviour change.
2026-10-04 21:44:03 +02:00
Dai Ha f6d1131d7a Merge remote-tracking branch 'origin/worker/726-unit2-75cb13-4' 2026-10-04 21:43:34 +02:00
Dai Ha cd1f04cbb4 Merge remote-tracking branch 'origin/worker/737-owner-key-ff061f-10'
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 53s
CI / build (push) Failing after 1m48s
2026-10-04 21:11:20 +02:00
Dai Ha efd9cdb983 fleetd #737 units 1+2: key tickets and turns on a stable owner, not a terminal
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Failing after 2m4s
Add Principal.ownerKey(): role-prefixed, keyed on name for a named lead and
a collaborator (survives a handover's terminal change), on terminal for a
worker, architect and observer, and on a distinct "anonymous" value for an
unauthenticated caller so that case no longer relies on Authz refusing it
first. The unnamed primary keeps a null key, preserving its primary-wide
ticket rule.

Thread that key through Task.creatorOwner, Rendezvous.Owner, poll,
pendingAsk and answer in place of a raw terminal, in both the MCP and REST
surfaces, so a named lead whose terminal changes can still poll and answer
its own delegations while a different lead is refused both.

Mutation evidence (each one-line change killed a named test, then reverted
to green):
- Principal.ownerKey() PRIMARY case made unconditional (dropped the
  null-name guard) -> ownerKeyCoversEveryRole dies:
  "expected: <null> but was: <leader:null>"
- OBSERVER case changed to use the "worker" prefix -> ownerKeyCoversEveryRole
  dies: "expected: <observer:term_observer> but was: <worker:term_observer>"
- prefixed() changed to drop the role prefix entirely -> both
  ownerKeyCoversEveryRole and rolePrefixesKeepLeadAndArchitectKeysDistinct
  die on a lead/architect key collision: "expected: <leader:opus> but was:
  <opus>"
- PRIMARY case changed to key on terminal instead of name ->
  rolePrefixesKeepLeadAndArchitectKeysDistinct and ownerKeyCoversEveryRole
  die: "expected: <leader:opus> but was: <leader:term_lead>"
- ARCHITECT case changed to key on the slot name instead of terminal ->
  same two tests die: "expected: <architect:opus> but was: <architect:design>"

All five mutations were caught by the existing test suite; no test needed
adding.
2026-10-04 20:42:44 +02:00
22 changed files with 724 additions and 249 deletions
+55 -53
View File
@@ -13,7 +13,7 @@ writes a handover file, and the new session reads that file and carries on.
- **By hand.** You write the file, then tell the operator where it is. The operator starts the new
session and points it at the file. This always works.
- **With `fleet_handover`** (fleetd #480, merged 2026-09-11). You ask fleetd to do the swap: it
checks the file, clears your pane, and tells the fresh session to read it. This needs
checks the file, restarts your pane, and tells the fresh session to read it. This needs
`leadRollover:` in `fleetd.yaml`; without it every action answers a clean refusal naming
`NOT_CONFIGURED`, and you fall back to the manual path. Section 11 below is the procedure.
@@ -159,38 +159,57 @@ fails.
`{action: "cancel", token}` drops a pending request without rolling.
**A change is coming: the roll will restart the process instead of sending `/clear` (fleetd #726
unit 2, written 2026-10-04).**
## 12. A roll restarts your process
Today a roll types `/clear` into your pane. Your `claude` process keeps running, so a newer CLI on
disk is never loaded. Unit 2 replaces that: the daemon ends the old pane, launches a fresh one,
waits for the new terminal to be recognised as a lead, and only then sends the bootstrap text.
The daemon ends your pane, launches a fresh one, waits for the new terminal to be recognised as a
lead, and only then sends the bootstrap text. Your `claude` process really exits, so a newer CLI on
disk is loaded. It does not type `/clear`.
**Everything below about `/clear` is accurate while the old jar is running.** Unit 2 was not merged
when this note was written, and a merge is not a deployment.
**The restart path is live in the code but has not run yet.** Measured 2026-10-05: the log holds no
`lead-rollover:` line since the current daemon started, and no occurrence of any new outcome name.
All 20 rolls recorded further down ran under the older `/clear` behaviour, so read them as history
rather than as evidence about your own roll. Re-measure with:
**How to tell which one is live: read your own tool list.** If `fleet_handover`'s description says
it will "clear your pane", the daemon is serving the old behaviour. If it names a restart, the new
behaviour is live. The description comes from the running daemon, so it cannot disagree with the
code that is actually loaded.
```bash
grep -c "lead-rollover: rolled" fleetd/fleetd.out # successful rolls
grep -c "lead-rollover:" fleetd/fleetd.out # positive control: must be larger
```
Two things change for you once it is live. The `status` outcomes are different: three new failures
replace the `/clear` ones. And the "never observed as WORKING after 8 consecutive IDLE/DONE polls"
warning described below can no longer appear, because that wait is deleted — so if you still see
it, the old jar is running. `TURN_NEVER_SETTLED` does not change, and still means nothing was
touched.
Run the control line too. A broken pattern returns a clean `0` that reads exactly like good news.
If the first number has grown past 20, somebody has rolled under the restart path, and this section
should be replaced with what they measured.
**Delete this note and rewrite the `/clear` paragraphs once the new jar is live.**
**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
seconds to die (`LeadRollover.PANE_DEATH_TIMEOUT_SECONDS`).
**Things that will surprise you:**
`fleet_handover{action: "status", token}` answers with one of these:
- **`accepted` does not mean your pane has been cleared.** It means every gate passed and the roll
| Outcome | What it means |
|---|---|
| `IN_PROGRESS` | still running; it always ends on one of the rows below |
| `ROLLED` | the roll succeeded |
| `TURN_NEVER_SETTLED` | your turn ran past `leadRollover.turnSettleSeconds`; nothing was touched |
| `OLD_PANE_NEVER_DIED` | your pane did not exit within the 10-second bound |
| `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 |
| `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
already gone, so you may never read them yourself — they are in the daemon log for whoever looks
next.
## 13. Things that will surprise you
- **`accepted` does not mean your pane has been restarted.** It means every gate passed and the roll
is scheduled to run once your current turn ends. Say your goodbye in the same turn — you will not
get another one.
- **If you are still running after that turn, the roll did not happen.** A roll that works clears
you, so surviving your own goodbye is itself the signal that it refused. Check with
- **If you are still running after that turn, the roll did not happen.** A roll that works ends your
process, so surviving your own goodbye is itself the signal that it refused. Check with
`fleet_handover{action: "status", token}`, using the token you confirmed. `TURN_NEVER_SETTLED`
means your turn ran past `leadRollover.turnSettleSeconds` and **no `/clear` was ever sent**: your
means your turn ran past `leadRollover.turnSettleSeconds` and **your pane was never ended**: your
context is intact and nothing was lost. Open a fresh request and retry. Never assume the roll
succeeded because `confirm` answered `accepted` — by the time it refuses, there is no caller left
to tell, so this check is the only thing that closes that gap.
@@ -216,46 +235,29 @@ touched.
**deferred, not hot** — it is read once at boot, so an edit does nothing until the daemon is
redeployed.
**Until fleetd #621 merges, the nudge text will tell you to ask the operator even where the
daemon no longer requires it.** `LeadHeartbeatLoop.contextNotice()` hardcodes "ask the operator"
and takes no config, so it cannot know. Trust the config value over the nudge text. Once #621 is
merged and deployed, the nudge matches the config and this warning can be deleted.
The context nudge tracks this value, so its text and the config agree (fleetd #621 —
`LeadHeartbeatLoop.contextNotice` takes `requireOperatorConfirm`). You still read the config
rather than the nudge, because the nudge only reaches you when your context is already high.
- **The roll can still refuse after `confirm` returns**, and by then there is no caller to tell.
Those outcomes are logged only, as `lead-rollover:` lines in the daemon log.
- **The bootstrap prompt works end to end. Measured 2026-09-22, re-measured 2026-10-04.** This used
to say the fix was unproven (fleetd #489) and told you to expect a failure. That is no longer
true. On 2026-09-22 the daemon log held four `lead-rollover: rolled` lines. On 2026-10-04 it holds
**20**, against a control of 86 `lead-rollover:` lines. Each roll cleared the old lead and started
a fresh session against the handover file, with the configured `bootstrapText` arriving as its
first message. No context was lost. The old `Unknown command: /clearFresh` failure from 2026-09-12
does not appear in the log at all.
- **The bootstrap prompt works end to end — measured under the older `/clear` path.** On 2026-09-22
the daemon log held four `lead-rollover: rolled` lines; on 2026-10-04 it held **20**, against a
control of 86 `lead-rollover:` lines. Each roll started a fresh session against the handover file,
with the configured `bootstrapText` arriving as its first message, and no context was lost. So the
bootstrap half of the roll is proven, and that half did not change. Section 12 says how to check
whether anything has rolled under the restart path since.
19 of the 20 carry an `elapsedMs`: median 16507 ms, maximum 48261 ms, and two above 45000 ms. That
figure times the **whole** roll, and the wait for your own turn to end dominates it, so do not
read it as the cost of the clear. Expect a roll to take tens of seconds, and do not treat a slow
one as a failed one. Re-measure all of these with:
```bash
grep -c "lead-rollover: rolled" fleetd/fleetd.out # successful rolls
grep -c "lead-rollover:" fleetd/fleetd.out # positive control: must be larger
grep -c "Unknown command" fleetd/fleetd.out # the old failure: expect 0
```
Run the control line too. A broken pattern returns a clean `0` that reads exactly like good news.
If the first number stops growing across rolls, or `Unknown command` returns anything above 0,
the bootstrap has regressed and this paragraph is stale again.
figure times the **whole** roll, and the wait for your own turn to end dominates it. Expect a roll
to take tens of seconds, and do not treat a slow one as a failed one. The restart path adds a pane
death and a relaunch to that work, so expect it to be slower rather than faster — but nobody has
measured it, so do not quote a number for it.
**You still write the file before you confirm, and never the other way round.** That order is not
about the bootstrap being unreliable. It is what the daemon checks: the handover file must have
been modified *after* the open request, or `confirm` refuses it as stale.
- **One warning in the log is normal and is not a failure.** Every one of the three rolls above also
logged `/clear on term_… was never observed as WORKING after 8 consecutive IDLE/DONE polls —
releasing rather than wedging the roll`. The daemon could not see the pane go WORKING after
`/clear`, so it released instead of hanging. The roll then succeeded anyway. That is the safe
branch behaving correctly. Do not report it as a broken roll.
## Writing style
Write in plain English. Use everyday words, one idea per sentence, and active voice. Keep every
@@ -904,6 +904,32 @@ public final class Fleetd {
throw (T) t;
}
/**
* The {@link PrimaryRegistry} lookup for "which terminal currently hosts the lead named
* {@code name}" — the inverse of {@code liveLeadTerminals} (terminal id → lead name), read live
* on every call so a lead discovered, rolled, or lost since the last call is reflected without
* a restart. Returns {@code null} when no currently recognised lead carries that name — a name
* that is not a lead at all (an architect slot, a collaborator), or a lead whose tab the scan
* cannot currently place (just rolled, off-host, non-herdr).
*
* @param liveLeadTerminals terminal id → lead name for every CURRENTLY recognised lead, normally
* the same {@code leads} supplier {@code main} already builds for
* {@code HerdrRouter}/{@link #leadSeatLookup}
*/
static Function<String, String> currentTerminalForName(Supplier<Map<String, String>> liveLeadTerminals) {
return name -> {
if (name == null) {
return null;
}
for (var entry : liveLeadTerminals.get().entrySet()) {
if (name.equals(entry.getValue())) {
return entry.getKey();
}
}
return null;
};
}
/**
* fleetd #480: construct the {@link LeadRollover} executor only when {@code leadRollover:} is
* present at startup — the same presence gate {@code leadHeartbeat:} uses just above this
@@ -393,7 +393,8 @@ final class FleetdAssembly {
ports.leadMailboxOpener());
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null;
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal);
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal,
Fleetd.currentTerminalForName(leads));
// CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs.
if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) {
log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes "
@@ -145,6 +145,25 @@ public record Principal(Role role, String terminal, long pid, String name) {
return terminal != null && terminal.equals(sessionId);
}
/**
* Stable identity used to own tickets and open turns. The unnamed primary has no owner key so
* it can use the message layer's primary-wide ticket access rule.
*/
public String ownerKey() {
return switch (role) {
case PRIMARY -> name == null ? null : prefixed("leader", name);
case WORKER -> prefixed("worker", terminal);
case ARCHITECT -> prefixed("architect", terminal);
case COLLABORATOR -> prefixed("collaborator", name);
case OBSERVER -> prefixed("observer", terminal);
case ANONYMOUS -> "anonymous";
};
}
private static String prefixed(String role, String identity) {
return role + ":" + identity;
}
/** Short, non-sensitive description for audit lines and error details. */
public String describe() {
return switch (role) {
@@ -248,9 +248,9 @@ 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_NOT_RECOGNISED}, 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.
* #RELAUNCH_NEVER_READY}, {@link #RELAUNCH_NOT_RECOGNISED}, 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.
*/
IN_PROGRESS,
/**
@@ -459,12 +459,14 @@ public final class FleetMcp {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_send", req.arguments()),
str(req.arguments(), "sessionId"));
if (denied != null) return denied;
String caller = callerTerminal(exchange);
Principal caller = principal(exchange);
String callerTerminal = caller.terminal();
String callerOwner = caller.ownerKey();
// CB-548: only a PRIMARY caller may claim the legacy singleton "primary" fallback.
// An architect delegates as its own pane but must never become the fallback that
// no-delegation inbox nudges target as if it were the primary (the per-target
// delegation map does not cure the singleton).
recordPrimarySingleton(primaryRegistry, caller, principal(exchange));
recordPrimarySingleton(primaryRegistry, callerTerminal, caller);
Map<String, Object> a = req.arguments();
String target = str(a, "sessionId");
String content = str(a, "content");
@@ -480,18 +482,24 @@ 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), caller);
return answer(messages, turnId, content, timeoutMs(a), callerOwner);
}
// 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
// session lock and queued delivery — via the accepted-delivery callback, never at
// request time. A concurrent sender that times out BUSY therefore cannot steal a
// live turn's reply routing without ever owning the turn.
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
// Only a lead's name is ever resolvable back to a current terminal (PrimaryRegistry
// only looks it up among currently recognised leads) — an architect or collaborator
// 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);
// 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(), caller);
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), callerOwner);
};
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
// is "is this caller a worker at all", and it can only ever reply as itself.
@@ -514,7 +522,7 @@ public final class FleetMcp {
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null);
if (denied != null) return denied;
return status(messages, str(req.arguments(), "sessionId"), callerTerminal(exchange));
return status(messages, str(req.arguments(), "sessionId"), principal(exchange).ownerKey());
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> pollHandler =
(exchange, req) -> {
@@ -524,7 +532,8 @@ public final class FleetMcp {
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
if (denied != null) return denied;
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
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).
// Acking removes a reply from the inbox, so it is a drain, not a read.
@@ -727,7 +736,7 @@ public final class FleetMcp {
*/
static void recordPrimarySingleton(PrimaryRegistry registry, String callerTerminal, Principal caller) {
if (caller != null && caller.isPrimary()) {
registry.record(callerTerminal);
registry.record(callerTerminal, caller.name());
}
}
@@ -914,7 +923,7 @@ public final class FleetMcp {
*/
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
Long timeoutMs, Runnable onAccepted, Set<String> profiles,
String callerTerminal) {
String callerOwner) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -924,7 +933,7 @@ public final class FleetMcp {
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
try {
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerTerminal), timeout);
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerOwner), timeout);
} catch (HerdrException e) {
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
}
@@ -934,15 +943,15 @@ public final class FleetMcp {
* {@code fleet_send} carrying a {@code turnId}: the primary's answer to a worker's
* {@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 callerTerminal} must match the turn's recorded owner or this is refused.
* {@code callerOwner} must match the turn's recorded owner or this is refused.
*/
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs,
String callerTerminal) {
String callerOwner) {
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, callerTerminal), timeout);
return formatReply(messages.answer(turnId, content, timeout, callerOwner), timeout);
}
/**
@@ -1018,13 +1027,12 @@ public final class FleetMcp {
}
/**
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
* result back to that same caller — see {@link MessageService#poll(String, String)}.
* As above, recording {@code creator}'s owner key so a later {@code fleet_poll{ticket}} only
* hands the result back to the same caller — see {@link MessageService#poll(String, String)}.
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles,
String creatorTerminal) {
Runnable onAccepted, Set<String> profiles,
Principal creator) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -1032,7 +1040,7 @@ public final class FleetMcp {
if (targetError != null) {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
String ticket = messages.sendAsync(sessionId, content, onAccepted, creator);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
@@ -1222,13 +1230,12 @@ public final class FleetMcp {
}
/**
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
* client-supplied value.
* As above, refusing a ticket lookup whose caller owner key differs from the key that created it
* — see {@link MessageService#poll(String, String)}. {@code callerOwner} comes from the calling
* connection's resolved principal, never a client-supplied value.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId, String callerTerminal) {
String target, String coordId, String callerOwner) {
if (!isBlank(coordId)) {
return pollHeldPeerMail(leadChannel, coordId);
}
@@ -1242,7 +1249,7 @@ public final class FleetMcp {
if (isBlank(ticket)) {
return error("ticket (or target) is required");
}
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
MessageService.TaskView v = messages.poll(ticket, callerOwner);
if (v == null) {
return error("unknown ticket: " + ticket + " (never issued, or expired)");
}
@@ -1351,18 +1358,17 @@ public final class FleetMcp {
* {@code fleet_status}: the live lifecycle status of a worker session, plus — when the worker
* is paused mid-turn in an async {@code fleet_ask} — the open question and how to answer it, so
* a lead on its normal poll cadence does not need the ticket to notice. The question, its
* {@code turnId} and its ticket id are shown only to the caller whose terminal created that
* delegation, or to a caller with no terminal at all (the unnamed primary); any other caller
* still sees the base status. {@code callerTerminal} is the CALLING session's terminal id,
* resolved by the MCP layer from the connection, never a client-supplied value.
* {@code turnId} and its ticket id are shown only to the caller whose owner key created that
* delegation, or to the unnamed primary; any other caller still sees the base status.
* {@code callerOwner} comes from the calling connection's resolved principal.
*/
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerTerminal) {
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerOwner) {
if (isBlank(sessionId)) {
return error("sessionId is required");
}
try {
String base = messages.status(sessionId).name().toLowerCase();
MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerTerminal);
MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerOwner);
if (ask == null) {
return text(base);
}
@@ -6,6 +6,7 @@ import org.slf4j.LoggerFactory;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
/**
* Single-slot, thread-safe registry for the primary's herdr {@code terminal_id}.
@@ -18,13 +19,25 @@ import java.util.concurrent.atomic.AtomicReference;
* <p>The push loop ({@code ReplyPushLoop}) uses {@link #isKnown()} to decide
* whether active nudging is possible; an empty registry means the primary is
* off-host or non-herdr and delivery falls back to pull.
*
* <p><strong>A learned terminal can go stale; a configured lead's name cannot.</strong> A lead that
* is rolled (a fresh pane replacing the old one) keeps its name but gets a new {@code terminal_id}.
* So every terminal this class learns — the singleton and each per-target delegation — is recorded
* together with the delegating lead's name, when the caller carries one. {@link
* #currentPrimaryTerminal()} and {@link #nudgeTargetFor(String)} resolve that name back to a
* terminal through the live {@code currentTerminalForName} lookup before falling back to the
* terminal that was actually recorded. A caller with no name (an unnamed primary, an architect, a
* collaborator — none of those are leads a lookup keyed on lead names can resolve) is tracked by
* terminal alone, exactly as before this indirection existed.
*/
public final class PrimaryRegistry {
private static final Logger log = LoggerFactory.getLogger(PrimaryRegistry.class);
private final AtomicReference<String> terminal = new AtomicReference<>();
private final AtomicReference<String> primaryName = new AtomicReference<>();
private final boolean pinned;
private final Function<String, String> currentTerminalForName;
/**
* CB-532: worker terminal → the lead that delegated to it. The single slot above answers "who is
@@ -33,12 +46,32 @@ public final class PrimaryRegistry {
* other lead's delegations. This map answers the question that actually matters — "who is
* waiting on THIS worker" — and is what lets {@code primary.terminal} be retired.
*/
private final ConcurrentHashMap<String, String> leadByTarget = new ConcurrentHashMap<>();
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) {
}
/**
* @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank = unpinned)
*/
public PrimaryRegistry(String pinnedTerminal) {
this(pinnedTerminal, name -> null);
}
/**
* As above, with a live {@code lead name → current terminal} lookup — normally the inverse of
* the same {@code terminal_id → lead name} supplier {@code CallerResolver} and the lead-tab
* scan already read. A lookup that cannot place a name (it is not a currently recognised lead,
* or no lookup is wired) returns {@code null}, and every resolution here falls back to the
* terminal that was actually recorded.
*
* @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank =
* unpinned)
* @param currentTerminalForName lead name → its current terminal, or {@code null} if that name
* is not a currently recognised lead
*/
public PrimaryRegistry(String pinnedTerminal, Function<String, String> currentTerminalForName) {
if (pinnedTerminal != null && !pinnedTerminal.isBlank()) {
this.terminal.set(pinnedTerminal);
this.pinned = true;
@@ -46,19 +79,31 @@ public final class PrimaryRegistry {
} else {
this.pinned = false;
}
this.currentTerminalForName = currentTerminalForName != null ? currentTerminalForName : name -> null;
}
/**
* Record a terminal_id. No-op when:
* Record a terminal_id, with no lead name. No-op when:
* <ul>
* <li>the registry is pinned (config override),
* <li>{@code terminalId} is {@code null} or blank (non-herdr caller).
* </ul>
*/
public void record(String terminalId) {
record(terminalId, null);
}
/**
* As {@link #record(String)}, additionally recording the caller's name — present for a
* configured lead, {@code null} for an unnamed primary. The name is what lets {@link
* #currentPrimaryTerminal()} keep nudging the same lead across a roll even though its terminal
* changed.
*/
public void record(String terminalId, String name) {
if (pinned) return;
if (terminalId == null || terminalId.isBlank()) return;
String prev = terminal.getAndSet(terminalId);
primaryName.set(blankToNull(name));
if (prev == null) {
log.debug("primary terminal learned: {}", terminalId);
} else if (!prev.equals(terminalId)) {
@@ -67,7 +112,8 @@ public final class PrimaryRegistry {
}
/**
* Record that {@code leadTerminal} owns the accepted delegation of worker {@code target} (CB-532).
* Record that {@code leadTerminal} owns the accepted delegation of worker {@code target}
* (CB-532), with no lead name.
*
* <p>Called from the {@code MessageService} accepted-delivery hook — only after a send has won
* the session's send lock and queued delivery — where both halves are known (CB-548). It is
@@ -77,10 +123,19 @@ public final class PrimaryRegistry {
* lead that most recently delegated to it, which is the one waiting.
*/
public void recordDelegation(String target, String leadTerminal) {
recordDelegation(target, leadTerminal, null);
}
/**
* 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.
*/
public void recordDelegation(String target, String leadTerminal, String leadName) {
if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) {
return;
}
leadByTarget.put(target, leadTerminal);
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName)));
}
/** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */
@@ -99,19 +154,59 @@ public final class PrimaryRegistry {
* recorded delegation there is no right answer, so this returns empty rather than guessing —
* delivery degrades to pull, which is exactly what the durable inbox is for, instead of
* interrupting the wrong lead with someone else's result.
*
* <p>A delegation recorded with a name is resolved to that lead's <em>current</em> terminal
* first — see {@link #currentTerminalForName} — so a lead that has since been rolled is still
* reachable here, not just the pane that delegated the work originally.
*/
public Optional<String> nudgeTargetFor(String target) {
String lead = target == null ? null : leadByTarget.get(target);
return lead != null ? Optional.of(lead) : Optional.ofNullable(terminal.get());
Delegation delegation = target == null ? null : leadByTarget.get(target);
if (delegation != null) {
return Optional.of(resolveCurrent(delegation.terminal(), delegation.name()));
}
return currentPrimaryTerminal();
}
/** The known primary terminal, or empty if not yet learned (and not pinned). */
/**
* 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
* nudge destination which survives a lead roll want {@link #currentPrimaryTerminal()} instead.
*/
public Optional<String> primaryTerminal() {
return Optional.ofNullable(terminal.get());
}
/**
* The terminal to nudge for the singleton primary right now: the recorded name resolved to its
* current terminal when one was recorded and is still a recognised lead, otherwise the terminal
* that was actually recorded — empty only when nothing has been learned or pinned at all.
*/
public Optional<String> currentPrimaryTerminal() {
String learned = terminal.get();
if (learned == null) {
return Optional.empty();
}
return Optional.of(resolveCurrent(learned, primaryName.get()));
}
/** {@code true} once a terminal has been recorded (or was pinned at construction). */
public boolean isKnown() {
return terminal.get() != null;
}
/**
* {@code learnedTerminal}, unless {@code name} is non-null and {@code currentTerminalForName}
* currently places that name at a different, live terminal — in which case the live one wins.
*/
private String resolveCurrent(String learnedTerminal, String name) {
if (name == null) {
return learnedTerminal;
}
String current = currentTerminalForName.apply(name);
return current != null && !current.isBlank() ? current : learnedTerminal;
}
private static String blankToNull(String s) {
return s == null || s.isBlank() ? null : s;
}
}
@@ -307,12 +307,12 @@ public final class LeadHeartbeatLoop {
* (mirroring {@link ReplyPushLoop#tick(String)}) so tests can drive it directly with a fake clock and a
* fake {@link AgentControl} instead of racing the scheduler thread. */
void tick() {
boolean leadKnown = primaryRegistry.primaryTerminal().isPresent();
boolean leadKnown = primaryRegistry.currentPrimaryTerminal().isPresent();
FleetState fleet = snapshot(inbox, roster);
AgentStatus status = AgentStatus.UNKNOWN;
LeadContextGauge.Reading reading = LeadContextGauge.Reading.unknown();
if (leadKnown) {
String leadTerminal = primaryRegistry.primaryTerminal().orElseThrow();
String leadTerminal = primaryRegistry.currentPrimaryTerminal().orElseThrow();
try {
status = agents.status(leadTerminal);
} catch (RuntimeException e) {
@@ -367,7 +367,7 @@ public final class LeadHeartbeatLoop {
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
// notice at all, defeating the very check this fixes.
String notice = contextNotice(contextHighNudge, reading, contextNotified, requireOperatorConfirm);
var lead = primaryRegistry.primaryTerminal();
var lead = primaryRegistry.currentPrimaryTerminal();
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
// actually included in the text, and the send reached the pane without throwing. Whenever no
@@ -1,5 +1,6 @@
package dev.ltms.fleet.msg;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrRouter;
@@ -282,18 +283,15 @@ public final class MessageService {
*/
private volatile boolean askTimedOut;
/**
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
* created through an overload that does not record one. {@link #poll(String, String)}
* compares a polling caller's own terminal against this field before handing back the
* ticket's state.
* The owner key of the caller whose {@code fleet_send{wait:false}} created this ticket, or
* {@code null} for the unnamed primary and overloads that do not record a caller.
*/
private final String creatorTerminal;
private final String creatorOwner;
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
private Task(String ticket, String target, LongSupplier nowNanos, String creatorOwner) {
this.ticket = ticket;
this.target = target;
this.creatorTerminal = creatorTerminal;
this.creatorOwner = creatorOwner;
this.createdNanos = nowNanos.getAsLong();
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
@@ -357,7 +355,7 @@ public final class MessageService {
private final AtomicLong ticketSeq = new AtomicLong();
/**
* Minted once per {@code MessageService} instance and folded into every ticket id (see
* {@link #sendAsync(String, String, Runnable, String)}). {@link #ticketSeq} alone restarts at
* {@link #sendAsync(String, String, Runnable, Principal)}). {@link #ticketSeq} alone restarts at
* zero for every instance, so without this a ticket id can be reused across instances and
* resolve to an unrelated {@link Task} with no error; this nonce makes that impossible, because
* an id minted by one instance can never match the id space of another.
@@ -938,13 +936,13 @@ public final class MessageService {
/**
* Deliver {@code content} to {@code target} (a herdr {@code terminal_id}) and block until the
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. {@code callerTerminal}
* is the terminal of the caller making this call — {@code null} for the unnamed primary — and is
* recorded as the turn's owner, the only caller {@link #answer(String, String, long, String)} will
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. {@code callerOwner}
* identifies the caller making this call and is recorded as the turn's owner. It is the only
* caller {@link #answer(String, String, long, String)} will
* later accept an answer from if the worker pauses mid-turn to ask.
*/
public Reply send(String target, String content, long timeoutMillis, String callerTerminal) {
return send(target, content, timeoutMillis, null, callerTerminal);
public Reply send(String target, String content, long timeoutMillis, String callerOwner) {
return send(target, content, timeoutMillis, null, callerOwner);
}
/**
@@ -959,13 +957,13 @@ public final class MessageService {
* acceptance means a concurrent sender that times out {@code BUSY} can never steal ownership it
* never earned. {@code null} disables the hook.
*/
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerTerminal) {
return send(target, content, timeoutMillis, onAccepted, null, callerTerminal);
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerOwner) {
return send(target, content, timeoutMillis, onAccepted, null, callerOwner);
}
/** Run a send, optionally stopping an async task that teardown already failed before acceptance. */
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task,
String callerTerminal) {
String callerOwner) {
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
@@ -985,7 +983,7 @@ public final class MessageService {
// open race). Opening first also means a throwing onAccepted (fired before enqueue) or an
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
// failed send leaves no stale waiter behind.
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target, Rendezvous.Owner.of(callerTerminal));
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target, Rendezvous.Owner.of(callerOwner));
// CB-640: this send now owns target's delivery, so any earlier stranded-reply or
// still-queued fact no longer describes the live state — clear both rather than let
// them outlive the send that supersedes them.
@@ -1165,20 +1163,20 @@ public final class MessageService {
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
* is unblocked so a reply that lands the instant it resumes is not lost.
*
* <p>{@code callerTerminal} is the terminal of the caller making this call — {@code null} for
* the unnamed primary. It is checked against the turn's recorded owner (the caller whose
* <p>{@code callerOwner} identifies the caller making this call. It is checked against the
* turn's recorded owner (the caller whose
* accepted delegation opened it, see {@link #send(String, String, long, String)} and
* {@link #sendAsync(String, String, Runnable, String)}) before anything else runs: a mismatch,
* {@link #sendAsync(String, String, Runnable, Principal)}) before anything else runs: a mismatch,
* including a turn with no owner on record at all, returns {@link Outcome#NOT_TURN_OWNER}
* without touching the rendezvous, the session lock, or any async task bookkeeping.
*/
public Reply answer(String turnId, String content, long timeoutMillis, String callerTerminal) {
public Reply answer(String turnId, String content, long timeoutMillis, String callerOwner) {
String workerSession = rendezvous.askSession(turnId);
if (workerSession == null) {
return new Reply(Outcome.STALE_TURN, null); // the ask lapsed (timed out or already answered)
}
Rendezvous.Owner owner = rendezvous.askOwner(turnId);
if (!Rendezvous.Owner.permits(owner, callerTerminal)) {
if (!Rendezvous.Owner.permits(owner, callerOwner)) {
return new Reply(Outcome.NOT_TURN_OWNER, null);
}
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
@@ -1321,16 +1319,16 @@ public final class MessageService {
}
/**
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
* differs from this one; {@code null} records no owner (a caller with no terminal — the
* unnamed primary — is always allowed to poll the result regardless).
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creator}'s owner key as this
* ticket's owner. The key is derived here from the resolved principal so callers cannot pass a
* terminal address where an owner identity is required.
*
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
public String sendAsync(String target, String content, Runnable onAccepted, Principal creator) {
String ticket = "task-" + ticketBootNonce + "-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
String creatorOwner = creator == null ? null : creator.ownerKey();
Task task = new Task(ticket, target, nowNanos, creatorOwner);
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
@@ -1361,7 +1359,7 @@ public final class MessageService {
}
asyncExecutor.submit(() -> {
try {
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorTerminal);
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorOwner);
if (result.outcome() == Outcome.QUESTION) {
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
// just after resolveQuestion wakes this thread.
@@ -1393,9 +1391,9 @@ public final class MessageService {
}
/**
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
* As {@link #poll(String, String)}, with no caller owner key — the unnamed primary's ticket rule
* checked, so this overload must only be used where the caller's identity is otherwise
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
* irrelevant.
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
@@ -1403,18 +1401,18 @@ public final class MessageService {
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
* Refuses a {@code callerOwner} that differs from the owner that created the ticket (see
* {@link #sendAsync(String, String, Runnable, Principal)}) with a {@link Phase#FAILED} view that
* carries no reply text. The unnamed primary has a {@code null} owner key and is never refused.
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
*/
public TaskView poll(String ticket, String callerTerminal) {
public TaskView poll(String ticket, String callerOwner) {
Task task = tasks.get(ticket);
if (task == null) {
return null;
}
if (!ownsTicket(task, callerTerminal)) {
if (!ownsTicket(task, callerOwner)) {
return new TaskView(ticket, Phase.FAILED, null, null,
"forbidden: this ticket was created by a different session", null);
}
@@ -1454,14 +1452,13 @@ public final class MessageService {
}
/**
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
* equal the terminal recorded on the task; a task with no recorded terminal matches no
* terminal-bearing caller.
* Whether {@code callerOwner} may read {@code task}'s state. A {@code null} caller key is the
* unnamed primary and may read every ticket. Other callers must match the task's owner key. This
* differs from {@link Rendezvous.Owner#permits}: a missing rendezvous owner is not an authenticated
* unnamed primary, so that gate refuses every caller when no owner was recorded.
*/
private static boolean ownsTicket(Task task, String callerTerminal) {
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
private static boolean ownsTicket(Task task, String callerOwner) {
return callerOwner == null || callerOwner.equals(task.creatorOwner);
}
/**
@@ -1787,13 +1784,13 @@ public final class MessageService {
* {@code fleet_status} uses this to show a pending question without the caller needing the
* ticket. {@code null} when the session has no open async question (including a session mid a
* <em>blocking</em> {@code fleet_ask}, which has no {@link Task} to look up — see
* {@link PendingAsk}), or when {@code callerTerminal} does not own the task the question
* {@link PendingAsk}), or when {@code callerOwner} does not own the task the question
* belongs to (see {@link #ownsTicket(Task, String)}).
*/
public PendingAsk pendingAsk(String workerSession, String callerTerminal) {
public PendingAsk pendingAsk(String workerSession, String callerOwner) {
for (Task task : tasks.values()) {
Reply q = task.question;
if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerTerminal)) {
if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerOwner)) {
return new PendingAsk(task.ticket, q.text(), q.turnId());
}
}
@@ -73,23 +73,22 @@ public final class Rendezvous {
/**
* The caller whose accepted delegation opened a turn — the only caller allowed to answer it.
* A {@code null} terminal means the unnamed primary, an authenticated caller with no pane.
* A {@code null} owner key means the unnamed primary.
*/
public record Owner(String terminal) {
public record Owner(String ownerKey) {
public static final Owner UNNAMED_PRIMARY = new Owner(null);
public static Owner of(String terminal) {
return terminal == null ? UNNAMED_PRIMARY : new Owner(terminal);
public static Owner of(String ownerKey) {
return ownerKey == null ? UNNAMED_PRIMARY : new Owner(ownerKey);
}
/**
* Whether {@code callerTerminal} matches {@code owner}. A {@code null} owner means no
* owner was ever recorded, and that state matches no caller, not even one whose own
* terminal is {@code null} — "no record" and "recorded as the unnamed primary" are
* different states.
* Whether {@code callerOwner} matches {@code owner}. A {@code null} owner means no owner was
* recorded, so it matches no caller. {@link #UNNAMED_PRIMARY} records the unnamed primary
* with an owner object whose key is {@code null}.
*/
public static boolean permits(Owner owner, String callerTerminal) {
return owner != null && java.util.Objects.equals(owner.terminal(), callerTerminal);
public static boolean permits(Owner owner, String callerOwner) {
return owner != null && java.util.Objects.equals(owner.ownerKey(), callerOwner);
}
}
@@ -439,6 +439,15 @@ public final class ReplyPushLoop {
* — timeout, transport error, a decode error — is treated as still live and the binding is left
* alone, because guessing wrong here is unrecoverable while guessing "live" merely costs one more
* retry on the next tick, which {@link #decide} already tolerates.
*
* <p><strong>The fallback is probed too.</strong> {@code PrimaryRegistry.nudgeTargetFor} already
* resolves a named delegator to its current terminal before this method ever sees it, which
* keeps a rolled lead's per-target binding live. What that resolution cannot fix is a caller
* that was never recorded with a name at all — an unnamed primary, or a lead whose tab the
* scanner cannot currently see — where the fallback it returns is still the raw terminal last
* learned from call traffic. This method returns that fallback only after the same liveness
* check, and gives up for this tick (an empty result, exactly like "no lead known at all") rather
* than hand a caller a second stale address un-probed.
*/
private Optional<String> resolveLiveLead(String target) {
Optional<String> lead = primaryRegistry.nudgeTargetFor(target);
@@ -448,7 +457,12 @@ public final class ReplyPushLoop {
log.debug("push: lead {} delegated to for {} is no longer live, forgetting the stale binding "
+ "and falling back", lead.get(), target);
primaryRegistry.forgetDelegation(target);
return primaryRegistry.nudgeTargetFor(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);
return Optional.empty();
}
/**
@@ -664,7 +664,7 @@ public final class FleetApp {
return;
}
Principal caller = ctx.attribute(CALLER);
String callerTerminal = caller == null ? null : caller.terminal();
String callerOwner = caller == null ? null : caller.ownerKey();
JsonNode body;
try {
body = mapper.readTree(ctx.body());
@@ -690,19 +690,19 @@ public final class FleetApp {
// 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, callerTerminal), timeout);
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerOwner), timeout);
return;
}
if (!wait) {
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
String ticket = messages.sendAsync(id, content, null, callerTerminal);
String ticket = messages.sendAsync(id, content, null, caller);
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
return;
}
try {
writeReply(ctx, id, messages.send(id, content, timeout, callerTerminal), timeout);
writeReply(ctx, id, messages.send(id, content, timeout, callerOwner), timeout);
} catch (HerdrException e) {
herdrError(ctx, e);
}
@@ -883,10 +883,10 @@ public final class FleetApp {
body.put("ready", deliverable.test(id));
// A worker paused mid-turn in an async fleet_ask is otherwise invisible to a status
// poll — surface the open question and how to answer it, same as fleet_poll's
// Phase.ASKING view, but only to the caller whose terminal created that delegation, or
// to a caller with no terminal at all (the unnamed primary).
// Phase.ASKING view, but only to the caller whose owner key created that delegation, or
// to the unnamed primary.
Principal caller = ctx.attribute(CALLER);
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.terminal());
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.ownerKey());
if (ask != null) {
body.put("question", ask.question());
body.put("turnId", ask.turnId());
@@ -904,7 +904,7 @@ public final class FleetApp {
return;
}
Principal caller = ctx.attribute(CALLER);
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.terminal());
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;
@@ -0,0 +1,31 @@
package dev.ltms.fleet.auth;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
class PrincipalTest {
@Test
void ownerKeyCoversEveryRole() {
assertEquals("leader:opus", Principal.leader("opus", "term_lead", 1).ownerKey());
assertNull(Principal.primary(2).ownerKey());
assertEquals("worker:term_worker", Principal.worker("term_worker", 3).ownerKey());
assertEquals("architect:term_arch", Principal.architect("opus", "term_arch", 4).ownerKey());
assertEquals("collaborator:ops", Principal.collaborator("ops", "term_collab", 5).ownerKey());
assertEquals("observer:term_observer", Principal.observer("term_observer", 6).ownerKey());
assertEquals("anonymous", Principal.anonymous().ownerKey());
}
@Test
void rolePrefixesKeepLeadAndArchitectKeysDistinct() {
String lead = Principal.leader("opus", "term_lead", 1).ownerKey();
String architect = Principal.architect("design", "opus", 2).ownerKey();
assertEquals("leader:opus", lead);
assertEquals("architect:opus", architect);
assertNotEquals(lead, architect);
}
}
@@ -504,12 +504,12 @@ class FleetMcpAuthzTest {
}
/**
* {@code fleet_poll{ticket}} must thread the calling connection's own terminal into
* {@code fleet_poll{ticket}} must thread the calling connection's owner key into
* {@link MessageService#poll(String, String)}, so a worker cannot read a ticket a different
* session created.
*/
@Test
void theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll() throws Exception {
void theFleetPollHandlerActuallyThreadsCallerOwnerIntoPoll() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("pollHandler =");
@@ -525,18 +525,18 @@ class FleetMcpAuthzTest {
"control failed: the scraped pollHandler block contains no poll(messages, ...) call "
+ "at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
"the fleet_poll handler must thread callerTerminal(exchange) into poll(...), not omit "
assertTrue(handlerBlock.contains("principal(exchange).ownerKey()"),
"the fleet_poll handler must thread principal(exchange).ownerKey() into poll(...), not omit "
+ "it or pass a literal null -- block: " + handlerBlock);
}
/**
* {@code fleet_status} must thread the calling connection's own terminal into
* {@code fleet_status} must thread the calling connection's owner key into
* {@link FleetMcp#status(MessageService, String, String)}, so a caller that did not create a
* worker's open delegation cannot read its pending question through the status handler either.
*/
@Test
void theFleetStatusHandlerActuallyThreadsCallerTerminalIntoStatus() throws Exception {
void theFleetStatusHandlerActuallyThreadsCallerOwnerIntoStatus() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("statusHandler =");
@@ -552,8 +552,8 @@ class FleetMcpAuthzTest {
"control failed: the scraped statusHandler block contains no status(messages, ...) "
+ "call at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
"the fleet_status handler must thread callerTerminal(exchange) into status(...), not "
assertTrue(handlerBlock.contains("principal(exchange).ownerKey()"),
"the fleet_status handler must thread principal(exchange).ownerKey() into status(...), not "
+ "omit it or pass a literal null -- block: " + handlerBlock);
}
@@ -329,12 +329,13 @@ class FleetMcpTest {
/**
* Same hijack and control through the fire-and-poll ({@code sendAsync}) path: the owner comes
* from the ticket's recorded creator terminal, not from a caller threaded through a live call.
* from the resolved principal, not from a caller argument threaded through a live call.
*/
@Test
void aDifferentCallersMcpAnswerIsRefusedForAnAsyncSendButTheRealOwnerSucceeds() throws Exception {
Principal owner = Principal.worker("term_owner", 1);
McpSchema.CallToolResult accepted =
FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), "term_owner");
FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), owner);
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
@@ -354,14 +355,15 @@ class FleetMcpTest {
assertEquals(MessageService.Phase.ASKING, asking.phase());
String turnId = asking.turnId();
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L,
"worker:term_attacker");
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, "term_owner"));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey()));
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -2011,13 +2013,13 @@ class FleetMcpTest {
/**
* {@code fleet_status}'s pending-ask block (the question, its {@code turnId} and its ticket)
* is shown only to the caller whose terminal created the delegation, or to a caller with no
* terminal at all (the unnamed primary) — a different terminal-bearing caller still sees the
* base status line, but none of the pending-ask fields.
* is shown only to the caller whose owner key created the delegation, or to the unnamed primary.
* A different caller still sees the base status line, but none of the pending-ask fields.
*/
@Test
void statusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, "term_creator");
void statusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
Principal creator = Principal.worker("term_creator", 1);
String ticket = messages.sendAsync(T, "task that asks", null, creator);
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -2035,7 +2037,7 @@ class FleetMcpTest {
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
String other = textOf(FleetMcp.status(messages, T, "term_other"));
String other = textOf(FleetMcp.status(messages, T, "worker:term_other"));
assertTrue(other.startsWith("idle"), "the base status must still be shown: " + other);
assertFalse(other.contains("which config file?"),
"a non-creating caller must not see the question text: " + other);
@@ -2044,10 +2046,12 @@ class FleetMcpTest {
assertFalse(other.contains(ticket),
"a non-creating caller must not see the ticket: " + other);
String creator = textOf(FleetMcp.status(messages, T, "term_creator"));
assertTrue(creator.contains("which config file?"), "the creator must see the question: " + creator);
assertTrue(creator.contains(asking.turnId()), "the creator must see the turnId: " + creator);
assertTrue(creator.contains(ticket), "the creator must see the ticket: " + creator);
String creatorStatus = textOf(FleetMcp.status(messages, T, creator.ownerKey()));
assertTrue(creatorStatus.contains("which config file?"),
"the creator must see the question: " + creatorStatus);
assertTrue(creatorStatus.contains(asking.turnId()),
"the creator must see the turnId: " + creatorStatus);
assertTrue(creatorStatus.contains(ticket), "the creator must see the ticket: " + creatorStatus);
String unnamed = textOf(FleetMcp.status(messages, T, null));
assertTrue(unnamed.contains("which config file?"),
@@ -2056,7 +2060,7 @@ class FleetMcpTest {
// Clean up the still-open ask so the background thread does not linger past the test.
String turnId = asking.turnId();
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(turnId, "config.yaml", 5000, "term_creator"));
() -> messages.answer(turnId, "config.yaml", 5000, creator.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -169,4 +169,101 @@ class PrimaryRegistryTest {
assertTrue(reg.nudgeTargetFor("term_worker").isEmpty());
assertTrue(reg.nudgeTargetFor(null).isEmpty());
}
// ── fleetd #737 unit 3: a named lead's terminal is resolved live, not just recorded ─────────
/**
* The whole point of carrying a name: a lead that has been rolled keeps its name but gets a
* fresh terminal. {@code currentTerminalForName} stands in for the live lead-tab scan here —
* it reports the lead now sits on a different terminal than the one that was recorded — and
* {@code nudgeTargetFor} must follow the name to that current terminal, not the stale one.
*/
@Test
void nudgeTargetForFollowsARolledLeadsNameToItsCurrentTerminal() {
var reg = new PrimaryRegistry(null, name -> "opus".equals(name) ? "term_opus_after_roll" : null);
reg.recordDelegation("term_worker", "term_opus_before_roll", "opus");
assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_worker").orElseThrow(),
"the name must be resolved to the lead's CURRENT terminal, not the one recorded "
+ "at delegation time");
}
/**
* The lookup cannot place every name — an architect/collaborator name (never a lead), or a lead
* whose tab the scan cannot currently see (just rolled, off-host, non-herdr). Either way the
* terminal actually recorded is still the right thing to try, exactly as before this unit.
*/
@Test
void nudgeTargetForFallsBackToTheRecordedTerminalWhenTheNameCannotBePlaced() {
var reg = new PrimaryRegistry(null, name -> null); // nothing is ever currently recognised
reg.recordDelegation("term_worker", "term_lead_recorded", "opus");
assertEquals("term_lead_recorded", reg.nudgeTargetFor("term_worker").orElseThrow());
}
/** The 2-arg {@code recordDelegation} overload records no name, so resolution never applies. */
@Test
void recordDelegationWithNoNameIsNeverResolvedByLookup() {
var reg = new PrimaryRegistry(null, name -> {
throw new AssertionError("a delegation recorded with no name must never consult the lookup");
});
reg.recordDelegation("term_worker", "term_lead");
assertEquals("term_lead", reg.nudgeTargetFor("term_worker").orElseThrow());
}
/** As {@link #nudgeTargetForFollowsARolledLeadsNameToItsCurrentTerminal}, for the singleton. */
@Test
void currentPrimaryTerminalFollowsARolledLeadsNameToItsCurrentTerminal() {
var reg = new PrimaryRegistry(null, name -> "sol".equals(name) ? "term_sol_after_roll" : null);
reg.record("term_sol_before_roll", "sol");
assertEquals("term_sol_after_roll", reg.currentPrimaryTerminal().orElseThrow());
assertEquals("term_sol_before_roll", reg.primaryTerminal().orElseThrow(),
"primaryTerminal() stays the raw recorded value — currentPrimaryTerminal() is the "
+ "one that resolves live");
}
@Test
void currentPrimaryTerminalFallsBackWhenTheNameCannotBePlaced() {
var reg = new PrimaryRegistry(null, name -> null);
reg.record("term_sol", "sol");
assertEquals("term_sol", reg.currentPrimaryTerminal().orElseThrow());
}
@Test
void currentPrimaryTerminalWithNoNameRecordedIsTheRawTerminal() {
var reg = new PrimaryRegistry(null, name -> {
throw new AssertionError("no name was ever recorded, the lookup must not be consulted");
});
reg.record("term_x");
assertEquals("term_x", reg.currentPrimaryTerminal().orElseThrow());
}
@Test
void currentPrimaryTerminalIsEmptyWhenNothingWasEverLearned() {
var reg = new PrimaryRegistry(null, name -> "anything");
assertTrue(reg.currentPrimaryTerminal().isEmpty());
}
/** A pin never carries a name, so a pinned registry's singleton resolution is always a no-op. */
@Test
void currentPrimaryTerminalForAPinIsNeverResolvedByLookup() {
var reg = new PrimaryRegistry("term_pinned", name -> {
throw new AssertionError("a pin carries no name, the lookup must not be consulted");
});
assertEquals("term_pinned", reg.currentPrimaryTerminal().orElseThrow());
}
/** {@code nudgeTargetFor}'s fallback to the singleton is the resolved one, not the raw one. */
@Test
void nudgeTargetForWithNoDelegationFallsBackToTheResolvedSingleton() {
var reg = new PrimaryRegistry(null, name -> "opus".equals(name) ? "term_opus_after_roll" : null);
reg.record("term_opus_before_roll", "opus");
assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_never_seen").orElseThrow());
}
}
@@ -648,6 +648,42 @@ class LeadHeartbeatLoopTest {
"the notice text must appear exactly once across all three sends: " + herdr.sentTexts());
}
// ── fleetd #737 unit 3: tick() nudges the lead's CURRENT terminal, not the learned one ───────
/**
* {@code tick()} reads {@code primaryRegistry.currentPrimaryTerminal()} both to check the lead's
* status and to send the nudge. Here the registry learned the lead's terminal under its name
* before a roll; {@code currentTerminalForName} stands in for the live lead-tab scan and reports
* the lead now sits on a different terminal. A correct tick must follow the name and nudge the
* new terminal — nudging the old one would mean the heartbeat lost the lead across its own roll.
*/
@Test
void tickNudgesTheLeadsCurrentTerminalAfterARoll() {
var herdr = new FailableHerdrClient("term_lead_after_roll");
var now = new AtomicLong(NOW);
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
inbox.own(WORKER);
inbox.publish(WORKER, "m1", "hello");
@SuppressWarnings("unchecked")
List<MemberSession>[] rosterBox = new List[]{List.of(new MemberSession("p1", WORKER, "prof",
MemberRole.DEV, "/cwd", null, 0, 0, 0, MemberSession.State.READY, null, null))};
AgentControl agents = new AgentControl(herdr);
PrimaryRegistry registry = new PrimaryRegistry(null,
name -> "opus".equals(name) ? "term_lead_after_roll" : null);
registry.record("term_lead_before_roll", "opus");
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, agents, inbox, scheduler, 5, 100_000);
LeadHeartbeatLoop loop = new LeadHeartbeatLoop(registry, agents, inbox, () -> rosterBox[0], pushLoop,
scheduler, now::get, IDLE_AFTER_NANOS, 100_000L, 0);
loop.tick(); // opens the idle window
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
loop.tick(); // past the quiet period, pending reply -> INJECT
assertEquals(List.of("term_lead_after_roll"), herdr.promptTargets(),
"the heartbeat must read and nudge the lead's CURRENT terminal, not the one learned "
+ "before the roll");
}
/**
* Fake herdr client for the four tests above: always reports {@code lead} as IDLE, records the
* {@code text} of every {@code agent.prompt} call, and can be told to throw on the very next
@@ -657,6 +693,7 @@ class LeadHeartbeatLoopTest {
private static final ObjectMapper MAPPER = new ObjectMapper();
private final String lead;
private final List<String> sentTexts = new ArrayList<>();
private final List<String> promptTargets = new ArrayList<>();
private boolean throwOnNextSend = false;
FailableHerdrClient(String lead) {
@@ -671,6 +708,11 @@ class LeadHeartbeatLoopTest {
return List.copyOf(sentTexts);
}
/** Every terminal an {@code agent.prompt} call named, in call order. */
List<String> promptTargets() {
return List.copyOf(promptTargets);
}
@Override
@SuppressWarnings("unchecked")
public JsonNode call(String method, Object params) {
@@ -681,12 +723,13 @@ class LeadHeartbeatLoopTest {
.put("agent_status", "idle"));
}
if ("agent.prompt".equals(method)) {
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
if (throwOnNextSend) {
throwOnNextSend = false;
throw new RuntimeException("simulated transient herdr send failure");
}
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
sentTexts.add(String.valueOf(p.get("text")));
promptTargets.add(String.valueOf(p.get("target")));
}
return MAPPER.createObjectNode();
}
@@ -19,7 +19,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
* {@link MessageService#poll(String)} overload. That overload skips the ownership check in
* {@code MessageService}'s {@code ownsTicket} entirely, so a caller of it can read any session's
* ticket. Every production caller must go through {@link MessageService#poll(String, String)}
* and pass a {@code callerTerminal} explicitly, even when it is {@code null}.
* and pass a {@code callerOwner} explicitly, even when it is {@code null}.
*
* <p>This reads each file's own source text rather than reflecting on compiled bytecode, because
* the risk is a future one-word edit at a call site, not a missing overload.
@@ -48,7 +48,7 @@ class MessageServicePollUsageTest {
+ "below proves nothing");
assertTrue(violations.isEmpty(), "found a call to the fail-open MessageService.poll(String) "
+ "overload, which skips the ownership check entirely -- pass a callerTerminal "
+ "overload, which skips the ownership check entirely -- pass a callerOwner "
+ "explicitly (even if null) through poll(String, String) instead: " + violations);
// CONTROL: the arity parser actually finds the two genuine two-argument call sites (the
@@ -58,7 +58,7 @@ class MessageServicePollUsageTest {
assertEquals(2, twoArgSites.size(), "control failed: expected exactly the two known "
+ "two-argument messages.poll(...) call sites, found: " + twoArgSites);
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetMcp.java")),
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerTerminal) "
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerOwner) "
+ "site among: " + twoArgSites);
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetApp.java")),
"control failed: did not find the FleetApp.java messages.poll(...) site among: "
@@ -4,6 +4,7 @@ import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
@@ -921,24 +922,26 @@ class MessageServiceTest {
assertNull(other.poll(ticket), "a ticket minted by a different instance must not resolve here");
}
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
// --- ticket ownership -----------------------------------------------------------------------
@Test
void pollByAnotherTerminalIsRefused() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_a");
void leadBIsRefusedFromLeadAsTicket() throws Exception {
Principal leadA = Principal.leader("opus", "term_a", 1);
Principal leadB = Principal.leader("sol", "term_b", 2);
String ticket = messages.sendAsync(T, "long task", null, leadA);
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "secret async result"), "a reply resolves the async send");
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, "term_a");
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, leadA.ownerKey());
assertNotNull(owner, "the creator must still be able to read its own ticket");
assertEquals(MessageService.Phase.DONE, owner.phase());
MessageService.TaskView refused = messages.poll(ticket, "term_b");
assertNotNull(refused, "a different terminal gets a refusal, not silence");
MessageService.TaskView refused = messages.poll(ticket, leadB.ownerKey());
assertNotNull(refused, "a different lead gets a refusal, not silence");
assertNotEquals(MessageService.Phase.DONE, refused.phase(),
"a different terminal must never see the ticket as DONE");
"a different lead must never see the ticket as DONE");
assertNull(refused.reply(), "a refusal must never carry the reply text");
assertFalse(String.valueOf(refused).contains("secret async result"),
"the reply text must not appear anywhere in the refused view");
@@ -946,14 +949,13 @@ class MessageServiceTest {
@Test
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
String ticket = messages.sendAsync(T, "long task", null,
Principal.leader("opus", "term_lead", 1));
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
// callerTerminal == null is the unnamed primary (resolved by token or loopback trust, with
// no herdr pane) — it must read a ticket a terminal-bearing lead created.
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
assertNotNull(view, "the unnamed primary must be able to read any ticket");
assertEquals(MessageService.Phase.DONE, view.phase());
@@ -961,19 +963,44 @@ class MessageServiceTest {
}
@Test
void creatorReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
void namedLeadCanPollItsTicketAfterItsTerminalChanges() throws Exception {
Principal oldLead = Principal.leader("opus", "term_OLD", 1);
Principal newLead = Principal.leader("opus", "term_NEW", 2);
assertNotEquals(oldLead.terminal(), newLead.terminal(), "the test requires different terminals");
String ticket = messages.sendAsync(T, "long task", null, oldLead);
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "own result"), "a reply resolves the async send");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, "term_creator");
assertNotNull(view, "the ticket's own creator must be able to read it");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, newLead.ownerKey());
assertNotNull(view, "the same named lead must read the ticket from its new terminal");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("own result", view.reply());
}
@Test
void anonymousOwnerKeyIsRefusedByTheTicketGateItself() {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "long task", null, lead);
MessageService.TaskView refused = messages.poll(ticket, Principal.anonymous().ownerKey());
assertNotNull(refused);
assertEquals(MessageService.Phase.FAILED, refused.phase());
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
}
@Test
void architectOwnershipUsesTerminalRatherThanSlot() {
Principal oldArchitect = Principal.architect("opus", "term_OLD", 1);
Principal newArchitect = Principal.architect("opus", "term_NEW", 2);
String ticket = messages.sendAsync(T, "long task", null, oldArchitect);
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket, oldArchitect.ownerKey()).phase());
assertEquals(MessageService.Phase.FAILED, messages.poll(ticket, newArchitect.ownerKey()).phase());
}
@Test
void pollReportsACompletedTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task");
@@ -1950,13 +1977,13 @@ class MessageServiceTest {
}
/**
* A caller's own terminal must match the terminal that created the delegation to see its
* pending question; a different terminal-bearing caller sees nothing, and a caller with no
* terminal at all (the unnamed primary) always sees it.
* A caller's owner key must match the key that created the delegation to see its pending
* question. The unnamed primary always sees it.
*/
@Test
void pendingAskGatesTheQuestionByTheDelegationsCreatorTerminal() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, "term_creator");
void pendingAskGatesTheQuestionByTheDelegationsCreatorOwner() throws Exception {
Principal creator = Principal.worker("term_creator", 1);
String ticket = messages.sendAsync(T, "task that asks", null, creator);
awaitWaiting();
injectDelivery();
@@ -1964,10 +1991,10 @@ class MessageServiceTest {
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertNull(messages.pendingAsk(T, "term_other"),
"a caller whose terminal did not create the delegation must not see the question");
assertNull(messages.pendingAsk(T, "worker:term_other"),
"a caller whose key did not create the delegation must not see the question");
MessageService.PendingAsk own = messages.pendingAsk(T, "term_creator");
MessageService.PendingAsk own = messages.pendingAsk(T, creator.ownerKey());
assertNotNull(own, "the creating caller must see its own open question");
assertEquals("which config file?", own.question());
@@ -1976,7 +2003,7 @@ class MessageServiceTest {
assertEquals("which config file?", unnamed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_creator"));
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1984,13 +2011,12 @@ class MessageServiceTest {
}
/**
* A task created with no recorded creator terminal (a short {@code sendAsync} overload) must
* not hand its open question to any caller that does have a terminal — only a caller with no
* terminal at all may still see it.
* A task created with no recorded owner (a short {@code sendAsync} overload) must not hand its
* open question to a caller with an owner key. Only the unnamed primary may still see it.
*/
@Test
void pendingAskDeniesATerminalBearingCallerWhenTheTaskRecordsNoCreator() throws Exception {
String ticket = messages.sendAsync(T, "task that asks"); // no creatorTerminal recorded
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
@@ -1998,8 +2024,8 @@ class MessageServiceTest {
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertNull(messages.pendingAsk(T, "term_someone"),
"a terminal-bearing caller must not see a question whose task records no creator");
assertNull(messages.pendingAsk(T, "worker:term_someone"),
"a caller with an owner key must not see a question whose task records no creator");
assertNotNull(messages.pendingAsk(T, null),
"the unnamed primary must still see it even with no recorded creator");
@@ -2055,7 +2081,8 @@ class MessageServiceTest {
*/
@Test
void aDifferentCallersAnswerIsRefusedForAnAsyncSendDelegationButTheRealOwnerSucceeds() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, "term_owner");
Principal owner = Principal.worker("term_owner", 1);
String ticket = messages.sendAsync(T, "task that asks", null, owner);
awaitWaiting();
injectDelivery();
@@ -2063,7 +2090,8 @@ class MessageServiceTest {
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Reply hijacked = messages.answer(asking.turnId(), "evil.yaml", 500, "term_attacker");
MessageService.Reply hijacked = messages.answer(asking.turnId(), "evil.yaml", 500,
"worker:term_attacker");
assertEquals(MessageService.Outcome.NOT_TURN_OWNER, hijacked.outcome(),
"a caller that did not create this delegation must be refused, not served");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
@@ -2071,7 +2099,38 @@ class MessageServiceTest {
"a refused answer must not advance the async ticket's phase");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_owner"));
() -> messages.answer(asking.turnId(), "config.yaml", 5000, owner.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
assertEquals("done", awaitTicketPhase(ticket, MessageService.Phase.DONE).reply());
}
@Test
void namedLeadCanSeeAndAnswerAnAskAfterItsTerminalChangesWhileLeadBIsRefused() throws Exception {
Principal oldLead = Principal.leader("opus", "term_OLD", 1);
Principal newLead = Principal.leader("opus", "term_NEW", 2);
Principal leadB = Principal.leader("sol", "term_SOL", 3);
assertNotEquals(oldLead.terminal(), newLead.terminal(), "the test requires different terminals");
String ticket = messages.sendAsync(T, "task that asks", null, oldLead);
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertNotNull(messages.pendingAsk(T, newLead.ownerKey()),
"the same lead at its new terminal must see the pending ask");
assertNull(messages.pendingAsk(T, leadB.ownerKey()),
"another lead must not see the pending ask");
assertEquals(MessageService.Outcome.NOT_TURN_OWNER,
messages.answer(asking.turnId(), "evil.yaml", 500, leadB.ownerKey()).outcome());
assertFalse(ask.isDone(), "another lead must not resolve the worker's ask");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, newLead.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -241,18 +241,18 @@ class RendezvousTest {
@Test
void ownerPermitsIsFailClosedOnARecordAndThreeStatesAreDistinct() {
assertFalse(Rendezvous.Owner.permits(null, null),
"no owner on record refuses even a caller with no terminal");
assertFalse(Rendezvous.Owner.permits(null, "term_a"),
"no owner on record refuses a terminal-bearing caller too");
"no owner on record refuses even the unnamed primary");
assertFalse(Rendezvous.Owner.permits(null, "worker:term_a"),
"no owner on record refuses a caller with an owner key too");
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, null),
"the unnamed primary owner matches a caller with no terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, "term_a"),
"the unnamed primary owner does not match a terminal-bearing caller");
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_a"),
"a named owner matches the same terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_b"),
"a named owner refuses a different terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), null),
"the recorded unnamed primary matches a caller with a null owner key");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, "worker:term_a"),
"the recorded unnamed primary does not match another owner key");
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), "worker:term_a"),
"an owner matches the same key");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), "worker:term_b"),
"an owner refuses a different key");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), null),
"a named owner refuses the unnamed primary");
}
}
@@ -304,6 +304,40 @@ class ReplyPushLoopTest {
+ "still live");
}
// --- fleetd #737 unit 3: the fallback nudgeTargetFor returns must be probed too --------------
/**
* {@code resolveLiveLead} forgets a dead per-target binding and asks {@code PrimaryRegistry}
* again for a fallback. That fallback can be dead too — here PRIMARY, the pinned singleton, is
* affirmatively gone alongside DEAD_LEAD. A correct {@code resolveLiveLead} probes it exactly
* like the first lead and gives up for this tick rather than trust it unchecked.
*
* <p>This is checked through {@code onReplyQueued}/{@code isActive} rather than a send count:
* {@link ReplyPushLoop#onReplyQueued} only registers pending work and starts a schedule once
* {@code resolveLiveLead} returns a present value — a dead fallback that was trusted unprobed
* would already make this true, synchronously, with no tick or send needed to observe it. Pairs
* with {@link #aStaleLeadBindingFallsBackToTheLiveLeadInsteadOfNudgingADeadTerminal} as the
* positive control: same stale-DEAD_LEAD setup, but there the fallback (PRIMARY) is live and the
* nudge does fire — proving this test's "nothing happens" result comes from the fallback being
* dead, not from the assertion being unable to observe a nudge at all.
*/
@Test
void aDoublyDeadFallbackIsNeverTrustedAndStartsNoSchedule() {
registry.recordDelegation(WORKER, DEAD_LEAD);
var rec = new AllDeadHerdrClient(Set.of(DEAD_LEAD, PRIMARY));
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(1, 50);
loop.onReplyQueued(WORKER);
assertFalse(loop.isActive(),
"both the per-target binding and the fallback are dead, so resolveLiveLead must "
+ "return empty and onReplyQueued must never register pending work or start "
+ "a schedule — an unprobed fallback would start one here");
assertEquals(0, rec.promptTargets().size(), "nobody live was found, so nothing was ever sent");
}
// --- nudge format --------------------------------------------------------------------------
@Test
@@ -1248,6 +1282,52 @@ class ReplyPushLoopTest {
}
}
/**
* Fake herdr client for fleetd #737 unit 3: every terminal named in {@code deadTargets} reports
* {@code agent_not_found} from {@code agent.get} — unlike {@link DeadLeadHerdrClient}, which can
* only make one terminal dead, this can make a per-target binding AND its fallback dead in the
* same test. {@code agent.prompt} is recorded unconditionally (no liveness check of its own),
* so a test can tell "resolveLiveLead probed and correctly found nobody live" (no prompt call)
* apart from "resolveLiveLead trusted a dead fallback and sent into it anyway" (a prompt call to
* a terminal this fake has already declared gone).
*/
private static final class AllDeadHerdrClient implements HerdrClient {
private final Set<String> deadTargets;
private final List<String> promptTargets = Collections.synchronizedList(new ArrayList<>());
AllDeadHerdrClient(Set<String> deadTargets) {
this.deadTargets = deadTargets;
}
@Override
@SuppressWarnings("unchecked")
public JsonNode call(String method, Object params) {
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
if ("agent.get".equals(method)) {
String target = String.valueOf(p.get("target"));
if (deadTargets.contains(target)) {
throw new HerdrException("no such agent: " + target, "agent_not_found", null);
}
return MAPPER.createObjectNode()
.set("agent", MAPPER.createObjectNode()
.put("terminal_id", target)
.put("agent_status", "idle"));
}
if ("agent.prompt".equals(method)) {
promptTargets.add(String.valueOf(p.get("target")));
}
return MAPPER.createObjectNode();
}
List<String> promptTargets() {
return List.copyOf(promptTargets);
}
@Override
public void close() {
}
}
/**
* Fake herdr client for fleetd #368 review: {@code flakyTarget}'s FIRST {@code agent.get} call
* fails with a transient, non-{@code agent_not_found} {@code HerdrException} — a transport-level
@@ -148,11 +148,11 @@ class FleetAppAuthTest {
/**
* {@code GET /tasks/{ticket}} must resolve its caller the same way {@code allow(...)} does
* and thread that terminal into {@link MessageService#poll(String, String)}, not the
* and thread that owner key into {@link MessageService#poll(String, String)}, not the
* no-check overload that ignores who is asking.
*/
@Test
void theTaskStatusRouteActuallyThreadsTheCallersTerminalIntoPoll() throws Exception {
void theTaskStatusRouteActuallyThreadsTheCallersOwnerKeyIntoPoll() throws Exception {
String source = Files.readString(REST_SOURCE);
int start = source.indexOf("private void taskStatus(Context ctx) {");
@@ -168,8 +168,8 @@ class FleetAppAuthTest {
"control failed: the scraped taskStatus block contains no messages.poll( call at all "
+ "-- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("caller.terminal()"),
"the taskStatus route must thread the resolved caller's terminal into messages.poll(...), "
assertTrue(handlerBlock.contains("caller.ownerKey()"),
"the taskStatus route must thread the resolved caller's owner key into messages.poll(...), "
+ "not the no-check overload -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
"the taskStatus route must resolve its caller the same way allow(...) does, not via a "
@@ -178,11 +178,11 @@ class FleetAppAuthTest {
/**
* {@code GET /sessions/{id}/status} must resolve its caller the same way {@code allow(...)}
* does and thread that terminal into {@link MessageService#pendingAsk(String, String)}, not
* does and thread that owner key into {@link MessageService#pendingAsk(String, String)}, not
* the no-check overload that ignores who is asking.
*/
@Test
void theSessionStatusRouteActuallyThreadsTheCallersTerminalIntoPendingAsk() throws Exception {
void theSessionStatusRouteActuallyThreadsTheCallersOwnerKeyIntoPendingAsk() throws Exception {
String source = Files.readString(REST_SOURCE);
int start = source.indexOf("private void sessionStatus(Context ctx) {");
@@ -198,8 +198,8 @@ class FleetAppAuthTest {
"control failed: the scraped sessionStatus block contains no messages.pendingAsk( "
+ "call at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("caller.terminal()"),
"the sessionStatus route must thread the resolved caller's terminal into "
assertTrue(handlerBlock.contains("caller.ownerKey()"),
"the sessionStatus route must thread the resolved caller's owner key into "
+ "messages.pendingAsk(...), not the no-check overload -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
"the sessionStatus route must resolve its caller the same way allow(...) does, not via "
@@ -224,7 +224,8 @@ class FleetAppAuthTest {
Javalin otherWorkerApp = startOnSharedService(messages, herdr, 9001L); // -> term_shell
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
try {
String ticket = messages.sendAsync("term_a", "long task", null, "term_a");
String ticket = messages.sendAsync("term_a", "long task", null,
Principal.worker("term_a", FakeHerdr.WORKER_PID));
HttpResponse<String> refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, refused.statusCode());
@@ -289,12 +290,12 @@ class FleetAppAuthTest {
/**
* {@code GET /sessions/{id}/status} shows a worker's pending {@code fleet_ask} question, its
* {@code turnId} and its ticket only to the caller whose terminal created that delegation, or
* to a caller with no terminal at all (the unnamed primary) — a different terminal-bearing
* caller still sees the base status line, but none of the pending-ask fields.
* {@code turnId} and its ticket only to the caller whose owner key created that delegation, or
* to the unnamed primary. A different caller still sees the base status line, but none of the
* pending-ask fields.
*/
@Test
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception {
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
@@ -306,7 +307,8 @@ class FleetAppAuthTest {
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
try {
ObjectMapper mapper = new ObjectMapper();
String ticket = messages.sendAsync("term_target", "task that asks", null, "term_a");
String ticket = messages.sendAsync("term_target", "task that asks", null,
Principal.worker("term_a", FakeHerdr.WORKER_PID));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -345,7 +347,7 @@ class FleetAppAuthTest {
// Clean up the still-open ask so the background thread does not linger past the test.
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(turnId, "config.yaml", 5000, "term_a"));
() -> messages.answer(turnId, "config.yaml", 5000, "worker:term_a"));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {