Compare commits

...

6 Commits

Author SHA1 Message Date
Dai Ha 6ab3a81af7 fleetd #737 unit 6: stop the unnamed primary sharing the internal bypass
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m59s
ownsTicket treated a null callerOwner as "read everything", conflating
the internal no-check test seam with a real unnamed primary's owner
key. Give them two different values: a named INTERNAL_NO_OWNER_CHECK
marker for the test-only bypass, and null as just another owner key
that must equal the ticket's recorded creatorOwner (including a
null-to-null match, so an unnamed primary still owns its own tickets).
Applies to both ownsTicket call sites, poll and pendingAsk.
2026-10-05 05:24:37 +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
14 changed files with 535 additions and 260 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
@@ -145,6 +145,26 @@ 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 —
* {@code null} — and that is matched against a ticket's recorded owner the same way any other
* key is: it owns a ticket another unnamed primary created, and nothing else.
*/
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,18 @@ 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);
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, callerTerminal);
// 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 +516,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 +526,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.
@@ -914,7 +917,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 +927,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 +937,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 +1021,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 +1034,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 +1224,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 +1243,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 +1352,18 @@ 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 matches the
* delegation's creator — the unnamed primary matches only a delegation another unnamed
* primary created; 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);
}
@@ -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;
@@ -11,6 +12,7 @@ import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
@@ -282,18 +284,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 +356,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 +937,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 +958,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 +984,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 +1164,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 +1320,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 +1360,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,28 +1392,30 @@ public final class MessageService {
}
/**
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
* 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).
* As {@link #poll(String, String)}, but bypasses the ownership check entirely via
* {@link #INTERNAL_NO_OWNER_CHECK}. No production code calls this overload — it exists for
* tests that only need the ticket's state and have no caller identity to pass.
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
return poll(ticket, INTERNAL_NO_OWNER_CHECK);
}
/**
* 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.
* 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.
* 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's owner key is {@code null}, matched the same way
* as any other key — it reads a ticket another unnamed primary created, and is refused on a
* ticket a named caller created. 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 +1455,27 @@ 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.
* Marker passed as {@code callerOwner} to bypass the ownership check entirely. No
* {@link Principal#ownerKey()} ever produces this value — every real key is either
* {@code null} (the unnamed primary) or prefixed with its role, such as {@code "worker:"} or
* {@code "leader:"}. {@link #poll(String)} passes it; {@link #pendingAsk} has no matching
* no-check overload, so this stays package-private for the test that drives the bypass
* directly.
*/
private static boolean ownsTicket(Task task, String callerTerminal) {
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
static final String INTERNAL_NO_OWNER_CHECK = "internal:no-owner-check";
/**
* Whether {@code callerOwner} may read {@code task}'s state. {@code callerOwner} is matched
* against the task's recorded owner key by equality, including a {@code null} match — the
* unnamed primary's owner key is {@code null}, so it owns a ticket another unnamed primary
* created and nothing else, the same rule every other role follows. The only caller that
* reads any ticket is {@link #INTERNAL_NO_OWNER_CHECK}. 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 callerOwner) {
return INTERNAL_NO_OWNER_CHECK.equals(callerOwner)
|| Objects.equals(callerOwner, task.creatorOwner);
}
/**
@@ -1787,13 +1801,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);
}
}
@@ -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,15 @@ 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. An unnamed primary is
* held to the same rule: its owner key is {@code null}, which here does not match the named
* worker that created the delegation, so it sees none of the pending-ask fields either — the
* same as any other non-creating caller.
*/
@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 +2039,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,19 +2048,26 @@ 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?"),
"a caller with no terminal (the unnamed primary) must see the question: " + unnamed);
assertTrue(unnamed.startsWith("idle"), "the base status must still be shown: " + unnamed);
assertFalse(unnamed.contains("which config file?"),
"an unnamed primary must not see a question on a delegation a named worker created: " + unnamed);
assertFalse(unnamed.contains(asking.turnId()),
"a non-creating unnamed primary must not see the turnId: " + unnamed);
assertFalse(unnamed.contains(ticket),
"a non-creating unnamed primary must not see the ticket: " + unnamed);
// 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) {
@@ -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,59 +922,128 @@ 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");
}
@Test
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
void unnamedPrimaryIsRefusedFromANamedLeadsTicket() throws Exception {
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");
MessageService.TaskView refused = messages.poll(ticket, Principal.primary(1).ownerKey());
assertNotNull(refused, "a different owner gets a refusal, not silence");
assertEquals(MessageService.Phase.FAILED, refused.phase());
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
assertNull(refused.reply(), "a refusal must never carry the reply text");
assertFalse(String.valueOf(refused).contains("primary-visible result"),
"the reply text must not appear anywhere in the refused view");
}
/**
* Positive control for {@link #unnamedPrimaryIsRefusedFromANamedLeadsTicket}: without this,
* that test would pass just as well if {@code poll} refused every caller.
*/
@Test
void unnamedPrimaryReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, Principal.primary(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");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, Principal.primary(2).ownerKey());
assertNotNull(view, "an unnamed primary must be able to read a ticket another unnamed primary created");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@Test
void creatorReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
void theOneArgPollOverloadBypassesOwnershipEntirely() throws Exception {
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");
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 2000;
while (view == null || view.phase() != MessageService.Phase.DONE) {
if (System.currentTimeMillis() >= deadline) break;
view = messages.poll(ticket); // the one-arg, no-check overload -- no caller owner key at all
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(view, "the internal bypass must read a ticket owned by a named lead");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@Test
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 +2020,15 @@ 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 is held to the same rule as everyone else: its key is
* {@code null}, which here does not match the named worker that created this delegation, so
* it is refused too.
*/
@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,19 +2036,18 @@ 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());
MessageService.PendingAsk unnamed = messages.pendingAsk(T, null);
assertNotNull(unnamed, "a caller with no terminal (the unnamed primary) must always see the question");
assertEquals("which config file?", unnamed.question());
assertNull(messages.pendingAsk(T, Principal.primary(1).ownerKey()),
"an unnamed primary must not see a question on a delegation a named worker created");
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 +2055,42 @@ 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.
* Positive control for {@link #pendingAskGatesTheQuestionByTheDelegationsCreatorOwner}:
* without this, that test's refusal would pass just as well if {@code pendingAsk} refused
* every caller. Here the delegation's creator is itself an unnamed primary (owner key
* {@code null}), so another unnamed primary's {@code null} key must still match it.
*/
@Test
void pendingAskDeniesATerminalBearingCallerWhenTheTaskRecordsNoCreator() throws Exception {
String ticket = messages.sendAsync(T, "task that asks"); // no creatorTerminal recorded
void unnamedPrimarySeesItsOwnDelegationsPendingQuestion() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, Principal.primary(1));
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.PendingAsk unnamed = messages.pendingAsk(T, Principal.primary(2).ownerKey());
assertNotNull(unnamed, "an unnamed primary must see the question on a delegation another unnamed primary created");
assertEquals("which config file?", unnamed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(unnamed.turnId(), "config.yaml", 5000, null));
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());
}
/**
* {@code pendingAsk} has no public no-check overload the way {@link MessageService#poll}
* does, so this drives {@link MessageService#INTERNAL_NO_OWNER_CHECK} directly — the only way
* to exercise the bypass for this method.
*/
@Test
void pendingAskInternalBypassSeesAnyDelegationsPendingQuestion() throws Exception {
Principal creator = Principal.worker("term_creator", 1);
String ticket = messages.sendAsync(T, "task that asks", null, creator);
awaitWaiting();
injectDelivery();
@@ -1998,8 +2098,34 @@ 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");
MessageService.PendingAsk bypassed = messages.pendingAsk(T, MessageService.INTERNAL_NO_OWNER_CHECK);
assertNotNull(bypassed, "the internal bypass must see a question on a delegation a named worker created");
assertEquals("which config file?", bypassed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.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());
}
/**
* 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");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
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 +2181,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 +2190,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 +2199,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");
}
}
@@ -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 "
@@ -207,14 +207,15 @@ class FleetAppAuthTest {
}
/**
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket,
* while the creating worker and the unnamed primary both still read it. The ticket is minted
* directly on the shared {@link MessageService}, the same way {@code MessageServiceTest}
* drives {@link MessageService#poll(String, String)}, so this exercises only the REST poll
* route's own handling of the ownership already recorded on the ticket.
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket, and
* refuses an unnamed primary just the same: a named worker's ticket is not anyone else's to
* read, caller rank included. The ticket is minted directly on the shared
* {@link MessageService}, the same way {@code MessageServiceTest} drives
* {@link MessageService#poll(String, String)}, so this exercises only the REST poll route's
* own handling of the ownership already recorded on the ticket.
*/
@Test
void restPollRefusesADifferentWorkerButAllowsTheCreatorAndTheUnnamedPrimary() throws Exception {
void restPollRefusesADifferentWorkerAndAnUnnamedPrimaryOnANamedWorkersTicket() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
@@ -224,7 +225,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());
@@ -240,8 +242,10 @@ class FleetAppAuthTest {
HttpResponse<String> primary = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, primary.statusCode());
assertFalse(primary.body().contains("forbidden"),
"the unnamed primary must read any ticket: " + primary.body());
assertTrue(primary.body().contains("forbidden"),
"an unnamed primary must not read a ticket a named worker created: " + primary.body());
assertFalse(primary.body().contains("\"reply\""),
"a refusal must never carry reply text: " + primary.body());
} finally {
creatorApp.stop();
otherWorkerApp.stop();
@@ -249,6 +253,33 @@ class FleetAppAuthTest {
}
}
/**
* Positive control for
* {@link #restPollRefusesADifferentWorkerAndAnUnnamedPrimaryOnANamedWorkersTicket}: without
* this, that test's refusal would pass just as well if the route refused every caller. Here
* the ticket's creator is itself an unnamed primary, so another unnamed primary reading it
* over REST must still succeed.
*/
@Test
void restPollAllowsAnUnnamedPrimaryItsOwnTicket() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, new Rendezvous());
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
try {
String ticket = messages.sendAsync("term_a", "long task", null, Principal.primary(FakeHerdr.WORKER_PID));
HttpResponse<String> own = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, own.statusCode());
assertFalse(own.body().contains("forbidden"),
"an unnamed primary must read a ticket another unnamed primary created: " + own.body());
} finally {
primaryApp.stop();
}
}
/**
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
* caller's own terminal on the ticket it returns, so that caller can still poll its own
@@ -289,12 +320,13 @@ 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. An
* unnamed primary is held to the same rule: its owner key is {@code null}, which here does not
* match the named worker that created the delegation, so it sees none of the pending-ask
* fields either — the same as any other non-creating caller.
*/
@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 +338,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);
@@ -319,7 +352,8 @@ class FleetAppAuthTest {
MessageService.TaskView asking;
deadline = System.currentTimeMillis() + 3000;
do {
asking = messages.poll(ticket, null);
// the no-check overload: this is test plumbing waiting for ASKING, not the gate under test
asking = messages.poll(ticket);
Thread.sleep(5);
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
@@ -340,12 +374,15 @@ class FleetAppAuthTest {
JsonNode primary = mapper.readTree(
send(primaryApp.port(), "GET", "/sessions/term_target/status", null, null).body());
assertEquals("which config file?", primary.get("question").asText(),
"a caller with no terminal (the unnamed primary) must see the question");
assertEquals("idle", primary.get("status").asText(), "the base status must still be shown");
assertFalse(primary.has("question"),
"an unnamed primary must not see a question on a delegation a named worker created: " + primary);
assertFalse(primary.has("turnId"), "a non-creating unnamed primary must not see the turnId: " + primary);
assertFalse(primary.has("ticket"), "a non-creating unnamed primary must not see the ticket: " + primary);
// 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) {
@@ -396,7 +433,8 @@ class FleetAppAuthTest {
MessageService.TaskView asking;
deadline = System.currentTimeMillis() + 3000;
do {
asking = messages.poll(ticket, null);
// the no-check overload: this is test plumbing waiting for ASKING, not the gate under test
asking = messages.poll(ticket);
Thread.sleep(5);
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());