Compare commits

...

4 Commits

Author SHA1 Message Date
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
13 changed files with 293 additions and 182 deletions
@@ -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,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,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);
}
@@ -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);
}
}
@@ -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) {
@@ -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");
}
}
@@ -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) {