From efd9cdb98316ff2ee916041eac3656afa898aec4 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 4 Oct 2026 20:42:44 +0200 Subject: [PATCH] fleetd #737 units 1+2: key tickets and turns on a stable owner, not a terminal 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: but was: " - OBSERVER case changed to use the "worker" prefix -> ownerKeyCoversEveryRole dies: "expected: but was: " - prefixed() changed to drop the role prefix entirely -> both ownerKeyCoversEveryRole and rolePrefixesKeepLeadAndArchitectKeysDistinct die on a lead/architect key collision: "expected: but was: " - PRIMARY case changed to key on terminal instead of name -> rolePrefixesKeepLeadAndArchitectKeysDistinct and ownerKeyCoversEveryRole die: "expected: but was: " - ARCHITECT case changed to key on the slot name instead of terminal -> same two tests die: "expected: but was: " All five mutations were caught by the existing test suite; no test needed adding. --- .../java/dev/ltms/fleet/auth/Principal.java | 19 +++ .../java/dev/ltms/fleet/mcp/FleetMcp.java | 60 ++++----- .../dev/ltms/fleet/msg/MessageService.java | 91 +++++++------ .../java/dev/ltms/fleet/msg/Rendezvous.java | 19 ++- .../java/dev/ltms/fleet/rest/FleetApp.java | 16 +-- .../dev/ltms/fleet/auth/PrincipalTest.java | 31 +++++ .../dev/ltms/fleet/mcp/FleetMcpAuthzTest.java | 16 +-- .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 34 ++--- .../msg/MessageServicePollUsageTest.java | 6 +- .../ltms/fleet/msg/MessageServiceTest.java | 123 +++++++++++++----- .../dev/ltms/fleet/msg/RendezvousTest.java | 22 ++-- .../dev/ltms/fleet/rest/FleetAppAuthTest.java | 32 ++--- 12 files changed, 290 insertions(+), 179 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/auth/PrincipalTest.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/auth/Principal.java b/fleetd/src/main/java/dev/ltms/fleet/auth/Principal.java index 35e6ec2c..1444fc04 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/auth/Principal.java +++ b/fleetd/src/main/java/dev/ltms/fleet/auth/Principal.java @@ -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) { diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 3e007513..f42469e6 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -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 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 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 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 profiles, - String creatorTerminal) { + Runnable onAccepted, Set 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); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index 140b86ba..d91e9847 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -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 reply = rendezvous.open(target, Rendezvous.Owner.of(callerTerminal)); + CompletableFuture 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 before the worker * is unblocked so a reply that lands the instant it resumes is not lost. * - *

{@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 + *

{@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 * blocking {@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()); } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/Rendezvous.java b/fleetd/src/main/java/dev/ltms/fleet/msg/Rendezvous.java index 89ff957c..f38a8992 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/Rendezvous.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/Rendezvous.java @@ -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); } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index 55de1f17..14f4a259 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -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; diff --git a/fleetd/src/test/java/dev/ltms/fleet/auth/PrincipalTest.java b/fleetd/src/test/java/dev/ltms/fleet/auth/PrincipalTest.java new file mode 100644 index 00000000..bbecc82f --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/auth/PrincipalTest.java @@ -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); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java index c15acd33..162c79a9 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java @@ -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); } diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 69ce68b6..8fb7fd7b 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -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 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 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) { diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServicePollUsageTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServicePollUsageTest.java index 3bf33ad6..6eb7414d 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServicePollUsageTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServicePollUsageTest.java @@ -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}. * *

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: " diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 5bc27321..9008bedf 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -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 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 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 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 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")); diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousTest.java index ba317046..a27c39bf 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousTest.java @@ -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"); } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java index 9a2fa44c..6d74d0ce 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java @@ -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 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 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) {