From dd18bd1f38ceae258c2fa71a307352f2472c9ba0 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 4 Oct 2026 09:46:20 +0200 Subject: [PATCH 1/2] fleetd #715: gate answer() on the caller that owns the turn Record the turn's owner on the forward rendezvous waiter (Rendezvous.Owner, a three-state record: no record / unnamed primary / named terminal). A fresh fleet_ask copies that owner onto the ask turn; a coalesced duplicate ask keeps the first owner. answer() compares the answering caller against the stored owner before taking the session lock or reopening the resumed waiter, and a mismatch returns the new NOT_TURN_OWNER outcome instead of STALE_TURN, with no rendezvous/task/question cleanup. send() and answer() both drop their no-caller overloads; every call site in MessageService, FleetMcp and FleetApp now threads an explicit caller terminal through. AuthzTest pins that the unnamed-primary null allowance is safe only because ANONYMOUS never reaches ANSWER. Covers all four combinations of capture input x adapter: blocking-send and async-send delegations, answered through both FleetMcp and FleetApp, each with a hijack attempt refused and the real owner's answer succeeding as a control. --- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 17 +- .../dev/ltms/fleet/msg/MessageService.java | 55 ++++-- .../java/dev/ltms/fleet/msg/Rendezvous.java | 68 ++++++- .../java/dev/ltms/fleet/rest/FleetApp.java | 22 ++- .../java/dev/ltms/fleet/auth/AuthzTest.java | 16 ++ .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 113 ++++++++++-- .../ltms/fleet/msg/MessageServiceTest.java | 168 ++++++++++++------ .../dev/ltms/fleet/msg/RendezvousTest.java | 76 ++++++++ .../dev/ltms/fleet/rest/FleetAppAuthTest.java | 150 +++++++++++++++- 9 files changed, 580 insertions(+), 105 deletions(-) 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 7e4c85d1..2182abc4 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -480,7 +480,7 @@ 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)); + return answer(messages, turnId, content, timeoutMs(a), caller); } // 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 @@ -491,7 +491,7 @@ public final class FleetMcp { // 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()); + : send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), caller); }; // 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. @@ -908,7 +908,8 @@ public final class FleetMcp { * a BUSY interloper never claims a turn it did not win. {@code null} disables recording. */ static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content, - Long timeoutMs, Runnable onAccepted, Set profiles) { + Long timeoutMs, Runnable onAccepted, Set profiles, + String callerTerminal) { if (isBlank(sessionId) || isBlank(content)) { return error("sessionId and content are required"); } @@ -918,7 +919,7 @@ public final class FleetMcp { } long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs); try { - return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout); + return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerTerminal), timeout); } catch (HerdrException e) { return error("herdr error contacting session " + sessionId + ": " + e.getMessage()); } @@ -928,13 +929,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. */ - static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs) { + static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs, + String callerTerminal) { 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), timeout); + return formatReply(messages.answer(turnId, content, timeout, callerTerminal), timeout); } /** @@ -983,6 +986,8 @@ public final class FleetMcp { + "\" and content set to your answer; the worker resumes the same turn."); case STALE_TURN -> error("that question is no longer open — it timed out or was already " + "answered (turnId stale)"); + case NOT_TURN_OWNER -> error("this turn belongs to a different delegation — only the caller " + + "that opened it may answer it"); case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker " + r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]"); // fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call, 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 da466c93..cfb78c01 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -87,7 +87,7 @@ public final class MessageService { /** * The worker paused mid-turn to ask the primary a question (CB-205); {@code text} is the * question and {@code turnId} correlates the answer. Not terminal — the primary answers with - * {@link #answer(String, String, long)} and the turn resumes. + * {@link #answer(String, String, long, String)} and the turn resumes. */ QUESTION, /** Timed out after the message was delivered — the worker is still working. */ @@ -120,10 +120,16 @@ public final class MessageService { /** Another send to this session was in flight for the whole window. */ BUSY, /** - * An answer ({@link #answer(String, String, long)}) referenced a {@code turnId} that is no - * longer open — the worker's {@code fleet_ask} already timed out or was answered. + * An answer ({@link #answer(String, String, long, String)}) referenced a {@code turnId} + * that is no longer open — the worker's {@code fleet_ask} already timed out or was answered. */ - STALE_TURN + STALE_TURN, + /** + * An answer ({@link #answer(String, String, long, String)}) named a {@code turnId} that is + * still open, but the answering caller is not the caller whose accepted delegation opened + * it. Distinct from {@link #STALE_TURN} so a refusal is never reported as a lapsed turn. + */ + NOT_TURN_OWNER } /** @@ -133,7 +139,7 @@ public final class MessageService { * {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION}, * else {@code null} * @param turnId correlation id for a {@link Outcome#QUESTION} (answered via - * {@link #answer(String, String, long)}), else {@code null} + * {@link #answer(String, String, long, String)}), else {@code null} */ public record Reply(Outcome outcome, String text, String turnId) { /** A reply with no correlation id (the common terminal outcomes). */ @@ -665,7 +671,7 @@ public final class MessageService { case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, TIMED_OUT_UNCONFIRMED, BUSY -> "timeout"; case WORKER_FAILED -> "failed"; case BACKEND_EXHAUSTED -> "backend_exhausted"; - case STALE_TURN, QUESTION -> null; // not a completed delegation + case STALE_TURN, QUESTION, NOT_TURN_OWNER -> null; // not a completed delegation }; } @@ -924,14 +930,17 @@ 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. + * 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 + * later accept an answer from if the worker pauses mid-turn to ask. */ - public Reply send(String target, String content, long timeoutMillis) { - return send(target, content, timeoutMillis, null); + public Reply send(String target, String content, long timeoutMillis, String callerTerminal) { + return send(target, content, timeoutMillis, null, callerTerminal); } /** - * As {@link #send(String, String, long)}, but with an accepted-delivery hook. + * As {@link #send(String, String, long, String)}, but with an accepted-delivery hook. * *

{@code onAccepted} is invoked exactly once, once this send has won {@code target}'s send * lock and so become the accepted target turn — it runs before delivery is @@ -942,12 +951,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) { - return send(target, content, timeoutMillis, onAccepted, null); + public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerTerminal) { + return send(target, content, timeoutMillis, onAccepted, null, callerTerminal); } /** 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) { + private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task, + String callerTerminal) { long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L; ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock()); @@ -967,7 +977,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); + CompletableFuture reply = rendezvous.open(target, Rendezvous.Owner.of(callerTerminal)); // 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. @@ -1146,19 +1156,30 @@ public final class MessageService { * mid-turn (already picked up), so the answer flows back through its own open {@code fleet_ask} * 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 + * accepted delegation opened it, see {@link #send(String, String, long, String)} and + * {@link #sendAsync(String, String, Runnable, String)}) 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) { + public Reply answer(String turnId, String content, long timeoutMillis, String callerTerminal) { 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)) { + return new Reply(Outcome.NOT_TURN_OWNER, null); + } long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L; ReentrantLock lock = sessionLocks.computeIfAbsent(workerSession, _ -> new ReentrantLock()); if (!tryLock(lock, remainingMillis(deadlineNanos))) { return new Reply(Outcome.BUSY, null); } try { - CompletableFuture reply = rendezvous.open(workerSession); + CompletableFuture reply = rendezvous.open(workerSession, owner); // fleetd #575: this try used to open below, AFTER the Task lookup/registration and the // STALE_TURN early return that follows it — so that return was covered only by a // hand-rolled copy of the finally's own cleanup pair, not the finally itself. Widening the @@ -1332,7 +1353,7 @@ public final class MessageService { } asyncExecutor.submit(() -> { try { - Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task); + Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorTerminal); if (result.outcome() == Outcome.QUESTION) { // Keep the accepted owner until answer() finishes it. markAsyncQuestion may run // just after resolveQuestion wakes this thread. 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 12e53dd0..d67a9eea 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/Rendezvous.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/Rendezvous.java @@ -60,7 +60,7 @@ public final class Rendezvous { } /** A worker's open mid-turn question: the worker session it belongs to and the answer future. */ - private record AskWaiter(String session, CompletableFuture answer) { + private record AskWaiter(String session, CompletableFuture answer, Owner owner) { } /** @@ -70,7 +70,33 @@ public final class Rendezvous { public record AskTicket(String turnId, CompletableFuture answer, boolean fresh) { } - private final ConcurrentHashMap> waiters = new ConcurrentHashMap<>(); + /** + * 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. + */ + public record Owner(String terminal) { + public static final Owner UNNAMED_PRIMARY = new Owner(null); + + public static Owner of(String terminal) { + return terminal == null ? UNNAMED_PRIMARY : new Owner(terminal); + } + + /** + * 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. + */ + public static boolean permits(Owner owner, String callerTerminal) { + return owner != null && java.util.Objects.equals(owner.terminal(), callerTerminal); + } + } + + /** A registered forward waiter together with the owner its delegation was opened under. */ + private record ForwardWaiter(Owner owner, CompletableFuture future) { + } + + private final ConcurrentHashMap waiters = new ConcurrentHashMap<>(); /** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */ private final ConcurrentHashMap asks = new ConcurrentHashMap<>(); @@ -89,8 +115,17 @@ public final class Rendezvous { * code a double open is impossible; this is a tripwire for the day that no longer holds. */ public CompletableFuture open(String session) { + return open(session, null); + } + + /** + * Same as {@link #open(String)}, additionally recording {@code owner} as the caller whose + * delegation opened this waiter. A {@code null} owner records no owner at all — the + * fail-closed default {@link Owner#permits} refuses to everyone. + */ + public CompletableFuture open(String session, Owner owner) { CompletableFuture waiter = new CompletableFuture<>(); - CompletableFuture existing = waiters.putIfAbsent(session, waiter); + ForwardWaiter existing = waiters.putIfAbsent(session, new ForwardWaiter(owner, waiter)); if (existing != null) { throw new IllegalStateException( "rendezvous double-open for session " + session + " — a waiter is already registered"); @@ -105,7 +140,13 @@ public final class Rendezvous { * successful {@code open} after a finished turn requires this close to have happened first). */ public void close(String session, CompletableFuture waiter) { - waiters.remove(session, waiter); + waiters.computeIfPresent(session, (s, w) -> w.future() == waiter ? null : w); + } + + /** The owner recorded for {@code session}'s open waiter, or {@code null} if none is open. */ + public Owner ownerOf(String session) { + ForwardWaiter w = waiters.get(session); + return w == null ? null : w.owner(); } /** Whether a send is currently awaiting a resolution for {@code session}. */ @@ -119,7 +160,8 @@ public final class Rendezvous { * send (see the CB-116 note above) rather than whichever send happens to be waiting when they fire. */ public CompletableFuture currentWaiter(String session) { - return waiters.get(session); + ForwardWaiter w = waiters.get(session); + return w == null ? null : w.future(); } /** @@ -147,7 +189,7 @@ public final class Rendezvous { String turnId = openAsksBySession.computeIfAbsent(session, _ -> { String newTurnId = session + "#" + askSeq.incrementAndGet(); CompletableFuture answer = new CompletableFuture<>(); - AskWaiter waiter = new AskWaiter(session, answer); + AskWaiter waiter = new AskWaiter(session, answer, ownerOf(session)); asks.put(newTurnId, waiter); minted[0] = waiter; return newTurnId; @@ -183,6 +225,16 @@ public final class Rendezvous { return w == null ? null : w.session(); } + /** + * The owner recorded for {@code turnId} when its ask turn was freshly opened — the caller + * whose delegation {@link #answerAsk} must match. {@code null} if {@code turnId} is unknown or + * lapsed, or if the ask opened with no forward waiter owner on record. + */ + public Owner askOwner(String turnId) { + AskWaiter w = asks.get(turnId); + return w == null ? null : w.owner(); + } + /** * Resolve a worker's blocked {@code fleet_ask} with the primary's {@code answer}, unblocking it * to resume its turn. @@ -245,7 +297,7 @@ public final class Rendezvous { } private boolean complete(String session, Resolution resolution) { - CompletableFuture waiter = waiters.get(session); - return waiter != null && waiter.complete(resolution); + ForwardWaiter waiter = waiters.get(session); + return waiter != null && waiter.future().complete(resolution); } } 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 93254acb..55de1f17 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -663,6 +663,8 @@ public final class FleetApp { if (!allow(ctx, routeAction("POST /sessions/{id}/message"), id)) { return; } + Principal caller = ctx.attribute(CALLER); + String callerTerminal = caller == null ? null : caller.terminal(); JsonNode body; try { body = mapper.readTree(ctx.body()); @@ -688,20 +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), timeout); + writeReply(ctx, id, messages.answer(turnId, content, timeout, callerTerminal), timeout); return; } if (!wait) { // Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}. - Principal caller = ctx.attribute(CALLER); - String ticket = messages.sendAsync(id, content, null, caller == null ? null : caller.terminal()); + String ticket = messages.sendAsync(id, content, null, callerTerminal); ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted")); return; } try { - writeReply(ctx, id, messages.send(id, content, timeout), timeout); + writeReply(ctx, id, messages.send(id, content, timeout, callerTerminal), timeout); } catch (HerdrException e) { herdrError(ctx, e); } @@ -720,6 +721,10 @@ public final class FleetApp { case STALE_TURN -> ctx.status(409).json(Map.of( "sessionId", id, "error", "stale_turn", "detail", "that question is no longer open (timed out or already answered)")); + case NOT_TURN_OWNER -> ctx.status(403).json(Map.of( + "sessionId", id, "error", "not_turn_owner", + "detail", "this turn belongs to a different delegation — only the caller that " + + "opened it may answer it")); case REPLIED, COMPLETED_UNREPLIED -> { // replySource distinguishes a structured fleet_reply from the CB-106 completion // fallback (a scrape of the worker's transcript when it finished without replying). @@ -734,9 +739,10 @@ public final class FleetApp { // silent fall-through. That is exactly the bug this ticket exists to fix: // `default -> "done"` used to sit here and would have told a REST caller the // delegation completed for TIMED_OUT_UNCONFIRMED, the one outcome where delivery - // is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION and STALE_TURN can never - // actually reach this inner switch — the outer switch above always dispatches - // them first — but they still need an arm to keep this switch exhaustive. + // is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN and + // NOT_TURN_OWNER can never actually reach this inner switch — the outer switch + // above always dispatches them first — but they still need an arm to keep this + // switch exhaustive. "status", switch (reply.outcome()) { case TIMED_OUT_WORKING -> "working"; case TIMED_OUT_QUEUED -> "queued"; @@ -746,7 +752,7 @@ public final class FleetApp { case BUSY -> "busy"; case WORKER_FAILED -> "failed"; case BACKEND_EXHAUSTED -> "backend_exhausted"; - case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable + case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN, NOT_TURN_OWNER -> "done"; // unreachable }, "detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED ? "no reply within " + timeout + "ms; delivery is unconfirmed — the " diff --git a/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java b/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java index 6abafd2c..a9c5e380 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java @@ -30,6 +30,22 @@ class AuthzTest { assertTrue(Authz.isUnauthenticated(null)); } + /** + * {@code MessageService.answer}'s turn-ownership check treats a caller with no terminal as + * matching a turn recorded for the unnamed primary. That rule only stays safe because an + * unauthenticated caller — whose terminal is also {@code null} — never reaches {@code answer} + * at all: {@link #anonymousIsAuthorizedForNothing} already covers every action including + * {@code ANSWER}, but this test names the exact coupling so a future change to either side + * cannot drift without turning this test red. + */ + @Test + void anAnonymousCallerIsRefusedAnswerSoItCanNeverBeMistakenForTheUnnamedPrimary() { + assertFalse(Authz.permits(ANON, ANSWER, null), + "an anonymous caller, whose terminal is also null, must never reach answer() — the " + + "turn-ownership check's null-terminal match for the unnamed primary owner " + + "relies on this gate refusing it first"); + } + @Test void orchestrationBelongsToThePrimaryAlone() { for (Authz.Action a : new Authz.Action[]{SPAWN, STOP, SEND, DRAIN}) { 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 74c49f0c..303b350f 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -77,7 +77,7 @@ class FleetMcpTest { private void assertSendRoundTrips(String target, Set profiles) throws Exception { CompletableFuture send = CompletableFuture.supplyAsync( - () -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles)); + () -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null)); long deadline = System.currentTimeMillis() + 3000; while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) { Thread.sleep(5); @@ -103,7 +103,7 @@ class FleetMcpTest { void sendThenReplyRoundTrips() throws Exception { // fleet_send blocks; fleet_reply resolves it with the worker's structured answer. CompletableFuture send = CompletableFuture.supplyAsync( - () -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of())); + () -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null)); // Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now // queues in the inbox if no waiter is open, which would break the round-trip). @@ -181,7 +181,7 @@ class FleetMcpTest { String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"')); CompletableFuture answer = CompletableFuture.supplyAsync( - () -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L)); + () -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null)); assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS))); deadline = System.currentTimeMillis() + 3000; @@ -279,7 +279,7 @@ class FleetMcpTest { assertEquals("late reply", messages.drainReplies("term_a").getFirst().content()); CompletableFuture answer = CompletableFuture.supplyAsync( - () -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L)); + () -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null)); assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS))); while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) { Thread.sleep(5); @@ -288,6 +288,89 @@ class FleetMcpTest { assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS))); } + // --- fleetd #715: fleet_send{turnId} is gated on the caller that owns the turn ------------- + + /** + * A blocking {@code fleet_send} from one caller opens the turn; a {@code fleet_send{turnId}} + * from a different caller is refused as an error, and the real owner's answer still succeeds. + */ + @Test + void aDifferentCallersMcpAnswerIsRefusedForABlockingSendButTheRealOwnerSucceeds() throws Exception { + CompletableFuture send = CompletableFuture.supplyAsync( + () -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner")); + long deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting(T)); + + CompletableFuture ask = CompletableFuture.supplyAsync( + () -> FleetMcp.ask(messages, T, "which config?", 5000L)); + McpSchema.CallToolResult question = send.get(5, TimeUnit.SECONDS); + assertTrue(textOf(question).contains("[question]"), textOf(question)); + String questionText = textOf(question); + String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length()); + String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"')); + + McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker"); + assertTrue(hijacked.isError(), "a caller that did not open this turn must get an error, not an answer"); + assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask"); + + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner")); + assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS))); + deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + FleetMcp.reply(messages, T, Role.WORKER, "done"); + assertEquals("done", textOf(answer.get(5, TimeUnit.SECONDS))); + } + + /** + * 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. + */ + @Test + void aDifferentCallersMcpAnswerIsRefusedForAnAsyncSendButTheRealOwnerSucceeds() throws Exception { + McpSchema.CallToolResult accepted = + FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), "term_owner"); + String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim(); + + long deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting(T)); + + CompletableFuture ask = CompletableFuture.supplyAsync( + () -> FleetMcp.ask(messages, T, "which config?", 5000L)); + MessageService.TaskView asking = messages.poll(ticket); + deadline = System.currentTimeMillis() + 3000; + while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + asking = messages.poll(ticket); + } + assertEquals(MessageService.Phase.ASKING, asking.phase()); + String turnId = asking.turnId(); + + McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "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")); + assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS))); + deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + FleetMcp.reply(messages, T, Role.WORKER, "done"); + assertEquals("done", textOf(answer.get(5, TimeUnit.SECONDS))); + } + @Test void pollUnknownTicketIsAnError() { McpSchema.CallToolResult res = FleetMcp.poll(messages, "task-999", null); @@ -297,7 +380,7 @@ class FleetMcpTest { @Test void sendTimesOutWithAWorkingNote() { - McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of()); + McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null); assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error"); assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res)); } @@ -313,7 +396,7 @@ class FleetMcpTest { void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception { herdr.agentSendFailsWith("send_failed"); CompletableFuture send = CompletableFuture.supplyAsync( - () -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of())); + () -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null)); long deadline = System.currentTimeMillis() + 2000; while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { //noinspection BusyWait @@ -333,13 +416,13 @@ class FleetMcpTest { @Test void sendRejectsMissingArgs() { - assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError()); - assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of()).isError()); + assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null).isError()); + assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null).isError()); } @Test void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() { - McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol")); + McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null); McpSchema.CallToolResult async = FleetMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol")); assertTrue(blocking.isError()); @@ -358,7 +441,7 @@ class FleetMcpTest { assertSendRoundTrips("term_live_member", profiles); // A herdr-owned pane outside the bridge roster cannot be classified at accept time. - McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles); + McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null); assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time"); } @@ -450,7 +533,7 @@ class FleetMcpTest { void askThenAnswerRoundTrips() throws Exception { // The primary delegates and blocks; wait until its waiter is open before the worker asks. CompletableFuture send = CompletableFuture.supplyAsync( - () -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of())); + () -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null)); long deadline = System.currentTimeMillis() + 3000; while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) { //noinspection BusyWait @@ -472,7 +555,7 @@ class FleetMcpTest { // The primary answers via fleet_send(turnId); this blocks again for the worker's reply. CompletableFuture answer = CompletableFuture.supplyAsync( - () -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L)); + () -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null)); // The worker's ask returns the answer — it resumes the same turn. assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS))); @@ -498,7 +581,7 @@ class FleetMcpTest { @Test void answerToAStaleTurnIsAnError() { - McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L); + McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null); assertTrue(res.isError()); assertTrue(textOf(res).contains("no longer open"), textOf(res)); } @@ -1901,7 +1984,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)); + () -> messages.answer(turnId, "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); deadline = System.currentTimeMillis() + 3000; while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { @@ -1958,7 +2041,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)); + () -> messages.answer(turnId, "config.yaml", 5000, "term_creator")); 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/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 912da690..9deab1a7 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -78,7 +78,7 @@ class MessageServiceTest { } private CompletableFuture sendAsync(String content, long timeoutMillis) { - return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis)); + return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis, null)); } private void awaitWaiting() throws InterruptedException { @@ -188,7 +188,7 @@ class MessageServiceTest { localRendezvous, new InMemoryReplyInbox()); String brief = "Implement the requested change. ".repeat(20); CompletableFuture send = - CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000)); + CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000, (String) null)); long deadline = System.currentTimeMillis() + 2000; while (!localRendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { Thread.sleep(5); @@ -341,7 +341,7 @@ class MessageServiceTest { // The primary answers via fleet_send(turnId); this blocks again for the worker's reply. CompletableFuture answer = - CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000)); + CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null)); // The worker's ask returns the answer — it resumes the same turn. MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS); @@ -377,7 +377,7 @@ class MessageServiceTest { // The primary answers that one turnId; both asks unblock with the same answer. CompletableFuture answer = - CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000)); + CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null)); MessageService.AskResult a1 = ask1.get(5, TimeUnit.SECONDS); MessageService.AskResult a2 = ask2.get(5, TimeUnit.SECONDS); @@ -472,7 +472,7 @@ class MessageServiceTest { "the duplicate's own timeout elapses first"); CompletableFuture answered = - CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500)); + CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500, null)); MessageService.AskResult a = fresh.get(5, TimeUnit.SECONDS); assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome(), @@ -507,7 +507,7 @@ class MessageServiceTest { CompletableFuture lateAnswer = new CompletableFuture<>(); messages.setAskTimeoutRaceHookForTest(() -> - lateAnswer.complete(messages.answer(turnId, "too late", 500))); + lateAnswer.complete(messages.answer(turnId, "too late", 500, null))); try { MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS); assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome()); @@ -539,7 +539,7 @@ class MessageServiceTest { @Test void answeringAnUnknownTurnIsStale() { - MessageService.Reply r = messages.answer(T + "#999", "too late", 500); + MessageService.Reply r = messages.answer(T + "#999", "too late", 500, null); assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(), "an answer to a turn that never existed (or already lapsed) is stale, not a hang"); } @@ -572,7 +572,7 @@ class MessageServiceTest { messages.setAnswerAskLapseRaceHookForTest(() -> rendezvous.answerAsk(turnId, "raced in first")); try { - MessageService.Reply r = messages.answer(turnId, "too late", 500); + MessageService.Reply r = messages.answer(turnId, "too late", 500, null); assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(), "an ask already answered by the race must be seen as lapsed, not double-delivered"); } finally { @@ -596,7 +596,7 @@ class MessageServiceTest { void sendTimesOutBeforeDeliveryIsQueuedNotWorking() { // Nothing ever delivers the message and nothing resolves the send, so the reply future // times out with delivery still incomplete — the message is still queued for the worker. - MessageService.Reply r = messages.send(T, "never delivered", 50); + MessageService.Reply r = messages.send(T, "never delivered", 50, null); assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome(), "an undelivered send that times out is still queued, not working"); assertNull(r.text()); @@ -605,7 +605,7 @@ class MessageServiceTest { @Test void sendTimesOutAfterDeliveryIsStillWorking() throws Exception { CompletableFuture send = - CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 300)); + CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 300, null)); awaitWaiting(); injector.onStatus(T, AgentStatus.IDLE); // deliver — the delivered future now completes injector.onStatus(T, AgentStatus.WORKING); // worker starts but never replies @@ -630,7 +630,7 @@ class MessageServiceTest { void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() { messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE)); try { - MessageService.Reply reply = messages.send(T, "race delivery", 50); + MessageService.Reply reply = messages.send(T, "race delivery", 50, null); assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(), "cancel reporting DELIVERED means the worker received the timed-out message"); @@ -652,7 +652,7 @@ class MessageServiceTest { void sendTimesOutWithAttemptedDeliveryReportsUnconfirmedNotQueued() throws Exception { herdr.agentSendFailsWith("send_failed"); CompletableFuture send = - CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150)); + CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150, null)); awaitWaiting(); injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt → ATTEMPTED @@ -679,7 +679,7 @@ class MessageServiceTest { // The primary answers, unblocking the worker; but the worker never sends the follow-up // fleet_reply, so the answering send rides out its short window as still-working. - MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200); + MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200, null); assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answer.outcome(), "an answered worker that never replies times out as still working"); @@ -715,7 +715,7 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.QUESTION, q.outcome()); CompletableFuture answer = - CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000)); + CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null)); ask.get(5, TimeUnit.SECONDS); // worker resumed with the answer awaitWaiting(); // the answering call has (re)opened its own forward waiter @@ -723,7 +723,7 @@ class MessageServiceTest { MessageService.Reply done = answer.get(5, TimeUnit.SECONDS); assertEquals(MessageService.Outcome.REPLIED, done.outcome()); - MessageService.Reply probe = messages.send(T, "probe after normal reply", 300); + MessageService.Reply probe = messages.send(T, "probe after normal reply", 300, null); assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(), "the session lock must be released after a normal REPLIED answer(), or this bounded " + "follow-up send would come back BUSY instead of timing out on its own work"); @@ -756,13 +756,13 @@ class MessageServiceTest { // The primary answers, unblocking the worker; the worker never sends its follow-up // fleet_reply, so the answering call rides out its short window as still-working. CompletableFuture answer = - CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200)); + CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200, null)); MessageService.Reply answered = answer.get(5, TimeUnit.SECONDS); assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answered.outcome(), "an answered worker that never replies times out as still working"); ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer - MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300); + MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300, null); assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(), "the session lock must be released after a TIMED_OUT_WORKING answer(), or this bounded " + "follow-up send would come back BUSY instead of timing out on its own work"); @@ -791,7 +791,7 @@ class MessageServiceTest { new java.util.concurrent.atomic.AtomicReference<>(); Thread answerer = new Thread(() -> { try { - messages.answer(turnId, "config.yaml", 5000); + messages.answer(turnId, "config.yaml", 5000, null); caught.set(new AssertionError("expected answer() to throw")); } catch (Throwable t) { caught.set(t); @@ -813,7 +813,7 @@ class MessageServiceTest { ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer - MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300); + MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300, null); assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(), "the session lock must be released when answer()'s reply future fails exceptionally, " + "or this bounded follow-up send would come back BUSY instead of timing out on its own work"); @@ -840,7 +840,7 @@ class MessageServiceTest { new java.util.concurrent.atomic.AtomicReference<>(); Thread answerer = new Thread(() -> { try { - messages.answer(turnId, "config.yaml", 5000); + messages.answer(turnId, "config.yaml", 5000, null); caught.set(new AssertionError("expected answer() to throw")); } catch (Throwable t) { caught.set(t); @@ -859,7 +859,7 @@ class MessageServiceTest { ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer - MessageService.Reply probe = messages.send(T, "probe after interruption", 300); + MessageService.Reply probe = messages.send(T, "probe after interruption", 300, null); assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(), "the session lock must be released when answer()'s wait is interrupted, or this bounded " + "follow-up send would come back BUSY instead of timing out on its own work"); @@ -967,11 +967,11 @@ class MessageServiceTest { @Test void concurrentSendToSameSessionWhileFirstHoldsItIsBusy() throws Exception { CompletableFuture first = - CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000)); + CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000, null)); awaitWaiting(); // the first send now holds the session lock, blocked on its reply // A second send to the SAME session cannot take the lock within its short window. - MessageService.Reply busy = messages.send(T, "second", 100); + MessageService.Reply busy = messages.send(T, "second", 100, null); assertEquals(MessageService.Outcome.BUSY, busy.outcome(), "a second send while another holds the session is busy, not a hang"); assertNull(busy.text()); @@ -1001,13 +1001,13 @@ class MessageServiceTest { PrimaryRegistry reg = new PrimaryRegistry(null); // L accepts a delegation to W: the send wins the lock and queues delivery → L is recorded. CompletableFuture first = CompletableFuture.supplyAsync( - () -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L))); + () -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L), null)); awaitWaiting(); assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "an accepted send owns the delegation"); // A attempts W while L holds it → BUSY (lock never taken) → its hook never fires. - MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A)); + MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A), null); assertEquals(MessageService.Outcome.BUSY, busy.outcome()); assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "a BUSY send must not steal the delegator ownership it never earned"); @@ -1028,7 +1028,7 @@ class MessageServiceTest { void anAcceptedSendAfterThePriorOwnerFinishesBecomesTheNewOwner() throws Exception { PrimaryRegistry reg = new PrimaryRegistry(null); CompletableFuture first = CompletableFuture.supplyAsync( - () -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L))); + () -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L), null)); awaitWaiting(); injector.onStatus(T, AgentStatus.IDLE); injector.onStatus(T, AgentStatus.WORKING); @@ -1039,7 +1039,7 @@ class MessageServiceTest { // L finished; A's later accepted send takes the delegation over. CompletableFuture second = CompletableFuture.supplyAsync( - () -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A))); + () -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A), null)); awaitWaiting(); assertEquals(LEAD_A, reg.nudgeTargetFor(T).orElseThrow(), "an accepted send after the owner finished becomes the new delegator"); @@ -1060,7 +1060,7 @@ class MessageServiceTest { void answeringAnAskDoesNotRewriteDelegatorOwnership() throws Exception { PrimaryRegistry reg = new PrimaryRegistry(null); CompletableFuture send = CompletableFuture.supplyAsync( - () -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L))); + () -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L), null)); awaitWaiting(); assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "L owns the delegation"); @@ -1075,7 +1075,7 @@ class MessageServiceTest { // L answers the ask on the same turn; the answer path must not touch ownership. CompletableFuture answer = - CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000)); + CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); // the answering send reopened its forward waiter assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), @@ -1095,7 +1095,7 @@ class MessageServiceTest { void aThrowingAcceptedHookLeavesNoStaleWaiterOrQueuedOrphan() { assertThrows(IllegalStateException.class, () -> messages.send(T, "doomed", 500, - () -> { throw new IllegalStateException("ownership hook failed"); }), + () -> { throw new IllegalStateException("ownership hook failed"); }, null), "a throwing ownership hook fails the send loudly"); assertFalse(rendezvous.isWaiting(T), "the failed send must not leave a stale rendezvous waiter"); @@ -1223,7 +1223,7 @@ class MessageServiceTest { @Test void abandonFailsASendThatIsStillWaitingOnAReleasedSession() throws Exception { CompletableFuture send = - CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000)); + CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000, null)); awaitWaiting(); assertTrue(messages.abandon(T, "session released"), "a live waiter is abandoned"); @@ -1243,7 +1243,7 @@ class MessageServiceTest { @Test void abandonDoesNotOverwriteAnAlreadyResolvedSend() throws Exception { CompletableFuture send = - CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000)); + CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000, null)); awaitWaiting(); assertTrue(rendezvous.resolve(T, "the real answer")); @@ -1374,7 +1374,7 @@ class MessageServiceTest { assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase()); CompletableFuture answer = CompletableFuture.supplyAsync( - () -> messages.answer(asking.turnId(), "config.yaml", 5000)); + () -> messages.answer(asking.turnId(), "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); assertTrue(rendezvous.resolve(T, "done")); @@ -1414,7 +1414,7 @@ class MessageServiceTest { // same turnId must see it as lapsed rather than resolving a question nobody is waiting on. assertEquals(MessageService.AskOutcome.TIMED_OUT, ask.get(5, TimeUnit.SECONDS).outcome()); assertEquals(MessageService.Outcome.STALE_TURN, - messages.answer(asking.turnId(), "config.yaml", 200).outcome()); + messages.answer(asking.turnId(), "config.yaml", 200, null).outcome()); } @Test @@ -1451,7 +1451,7 @@ class MessageServiceTest { // The primary answers, but its own bounded wait for the worker's resumed turn is short and // expires before the worker (still genuinely working) gets back to it. - MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150); + MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150, null); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(), "the primary's own bounded wait gives up before the worker finishes resuming"); @@ -1480,7 +1480,7 @@ class MessageServiceTest { CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); - MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150); + MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150, null); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome()); @@ -1534,7 +1534,7 @@ class MessageServiceTest { messages.setFinishAsyncTaskRaceHookForTest(() -> messages.forgetTurnForTest(turnId)); CompletableFuture answer = - CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000)); + CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn @@ -1640,7 +1640,7 @@ class MessageServiceTest { String turnId = asking.turnId(); CompletableFuture answer = - CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000)); + CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer(), "the worker's own ask() call must have already unblocked with the primary's answer " + "before we force the race below"); @@ -1687,7 +1687,7 @@ class MessageServiceTest { MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); String turnId = asking.turnId(); - MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150); + MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150, null); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(), "the primary's own bounded wait must give up first, leaving turnId stamped with no " @@ -1753,7 +1753,7 @@ class MessageServiceTest { MessageService.TaskView asking = messages.poll(first); CompletableFuture answer = CompletableFuture.supplyAsync( - () -> messages.answer(asking.turnId(), "config.yaml", 5000)); + () -> messages.answer(asking.turnId(), "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); assertTrue(rendezvous.resolve(T, "done")); @@ -1787,7 +1787,7 @@ class MessageServiceTest { // The primary answers it — answer() resumes the turn and blocks for what comes next. CompletableFuture answer1 = CompletableFuture.supplyAsync( - () -> messages.answer(asking1.turnId(), "a1", 5000)); + () -> messages.answer(asking1.turnId(), "a1", 5000, null)); assertEquals("a1", ask1.get(5, TimeUnit.SECONDS).answer()); // Still in the SAME resumed turn — before replying — the worker asks again. @@ -1808,7 +1808,7 @@ class MessageServiceTest { // The primary answers the second question; the worker finally sends its real fleet_reply. CompletableFuture answer2 = CompletableFuture.supplyAsync( - () -> messages.answer(turnId2, "a2", 5000)); + () -> messages.answer(turnId2, "a2", 5000, null)); assertEquals("a2", ask2.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); assertTrue(rendezvous.resolve(T, "done")); @@ -1908,7 +1908,7 @@ class MessageServiceTest { assertEquals(asking.turnId(), pending.turnId()); CompletableFuture answer = CompletableFuture.supplyAsync( - () -> messages.answer(asking.turnId(), "config.yaml", 5000)); + () -> messages.answer(asking.turnId(), "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); assertNull(messages.pendingAsk(T, null), "an answered question is no longer pending"); awaitWaiting(); @@ -1943,7 +1943,7 @@ class MessageServiceTest { assertEquals("which config file?", unnamed.question()); CompletableFuture answer = CompletableFuture.supplyAsync( - () -> messages.answer(asking.turnId(), "config.yaml", 5000)); + () -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_creator")); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); assertTrue(rendezvous.resolve(T, "done")); @@ -1971,13 +1971,81 @@ class MessageServiceTest { "the unnamed primary must still see it even with no recorded creator"); CompletableFuture answer = CompletableFuture.supplyAsync( - () -> messages.answer(asking.turnId(), "config.yaml", 5000)); + () -> messages.answer(asking.turnId(), "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); assertTrue(rendezvous.resolve(T, "done")); assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome()); } + // --- fleetd #715: answer() is gated on the caller that owns the turn ----------------------- + + /** + * A blocking {@code fleet_send} from one caller opens the turn; a different caller's answer is + * refused with no side effect on the rendezvous or the worker's blocked {@code fleet_ask} — + * only the real owner can answer it. + */ + @Test + void aDifferentCallersAnswerIsRefusedForABlockingSendDelegationButTheRealOwnerSucceeds() throws Exception { + CompletableFuture send = CompletableFuture.supplyAsync( + () -> messages.send(T, "do X", 5000, "term_owner")); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); + MessageService.Reply q = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.QUESTION, q.outcome()); + assertFalse(rendezvous.isWaiting(T), "the forward waiter closes once the question surfaces"); + + // The hijack: a different caller answers the SAME turnId. + MessageService.Reply hijacked = messages.answer(q.turnId(), "evil.yaml", 500, "term_attacker"); + assertEquals(MessageService.Outcome.NOT_TURN_OWNER, hijacked.outcome(), + "a caller that did not open this turn must be refused, not served"); + assertFalse(rendezvous.isWaiting(T), "a refused answer must not open a forward waiter"); + assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask"); + assertEquals(T, rendezvous.askSession(q.turnId()), "a refused answer must leave the ask turn open"); + + // The control: the real owner answers the same turnId and the turn resumes normally. + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> messages.answer(q.turnId(), "config.yaml", 5000, "term_owner")); + 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()); + } + + /** + * Same hijack and control as the blocking case, but the delegation is opened through the + * fire-and-poll path ({@code sendAsync}) — the owner comes from the ticket's recorded creator + * terminal, not a caller argument threaded through a live blocking call. + */ + @Test + void aDifferentCallersAnswerIsRefusedForAnAsyncSendDelegationButTheRealOwnerSucceeds() throws Exception { + String ticket = messages.sendAsync(T, "task that asks", null, "term_owner"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = + 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"); + 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"); + assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(), + "a refused answer must not advance the async ticket's phase"); + + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_owner")); + 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()); + } + // --- CB-588: async ticket terminal nudges --------------------------------------------------- // // MessageService.reply's rendezvous fast path is exactly what an async ticket always takes @@ -2234,7 +2302,7 @@ class MessageServiceTest { // Clean up the still-open ask so the background thread does not linger past the test. CompletableFuture answer = CompletableFuture.supplyAsync( - () -> wiring.service().answer(asking.turnId(), "config.yaml", 5000)); + () -> wiring.service().answer(asking.turnId(), "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); assertTrue(rendezvous.resolve(T, "done")); @@ -2260,7 +2328,7 @@ class MessageServiceTest { .filter(c -> c.method().equals("agent.prompt")).count(); CompletableFuture answer = CompletableFuture.supplyAsync( - () -> wiring.service().answer(asking.turnId(), "config.yaml", 5000)); + () -> wiring.service().answer(asking.turnId(), "config.yaml", 5000, null)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1); @@ -2653,7 +2721,7 @@ class MessageServiceTest { void hasQueuedDeliveryIsTrueAfterAnUndeliveredSendTimesOut() { // Nothing ever delivers the message (never goes IDLE/BLOCKED), so the send times out with // TIMED_OUT_QUEUED — same setup as sendTimesOutBeforeDeliveryIsQueuedNotWorking above. - MessageService.Reply r = messages.send(T, "never delivered", 50); + MessageService.Reply r = messages.send(T, "never delivered", 50, null); assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome()); assertTrue(messages.hasQueuedDelivery(T), @@ -2665,7 +2733,7 @@ class MessageServiceTest { @Test void hasQueuedDeliveryClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception { - assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome()); + assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50, null).outcome()); assertTrue(messages.hasQueuedDelivery(T)); // A fresh send accepts delivery (opens its own waiter) — the stale queued fact is cleared. @@ -2682,7 +2750,7 @@ class MessageServiceTest { @Test void hasQueuedDeliveryClearsOnAbandon() { - assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome()); + assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50, null).outcome()); assertTrue(messages.hasQueuedDelivery(T)); messages.abandon(T, "session released"); 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 f659b31f..aafb3821 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousTest.java @@ -132,4 +132,80 @@ class RendezvousTest { assertNull(rendezvous.askSession(t.turnId()), "a closed ask is forgotten"); assertFalse(rendezvous.answerAsk(t.turnId(), "late"), "a closed ask can no longer be answered"); } + + // ── fleetd #715: turn ownership ──────────────────────────────────────────────────────────── + + @Test + void openWithNoOwnerRecordsNoOwnerAndOpenWithAnOwnerRecordsIt() { + rendezvous.open(W); + assertNull(rendezvous.ownerOf(W), "the no-owner overload records no owner at all"); + rendezvous.close(W, rendezvous.currentWaiter(W)); + + rendezvous.open(W, Rendezvous.Owner.of("term_lead")); + assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.ownerOf(W)); + } + + @Test + void ownerOfIsNullWhenNoWaiterIsOpen() { + assertNull(rendezvous.ownerOf(W), "no waiter open means no owner to report"); + } + + @Test + void openAskStampsTheForwardWaitersOwnerOntoTheFreshTurnOnly() { + rendezvous.open(W, Rendezvous.Owner.of("term_lead")); + + Rendezvous.AskTicket fresh = rendezvous.openAsk(W); + assertTrue(fresh.fresh()); + assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.askOwner(fresh.turnId()), + "a freshly-opened ask copies the forward waiter's current owner"); + + Rendezvous.AskTicket coalesced = rendezvous.openAsk(W); + assertFalse(coalesced.fresh()); + assertEquals(fresh.turnId(), coalesced.turnId()); + assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.askOwner(coalesced.turnId()), + "a coalesced duplicate ask rides the fresh owner's turn, unchanged"); + } + + @Test + void aSecondAskAfterTheFirstClosesStampsWhateverOwnerIsOpenAtThatLaterMoment() { + rendezvous.open(W, Rendezvous.Owner.of("term_lead")); + Rendezvous.AskTicket first = rendezvous.openAsk(W); + rendezvous.closeAsk(first.turnId()); + + // The forward waiter is reopened under a different owner before the second ask — mirrors + // answer() reopening with the owner it already checked, which can differ turn to turn. + rendezvous.close(W, rendezvous.currentWaiter(W)); + rendezvous.open(W, Rendezvous.Owner.of("term_other")); + + Rendezvous.AskTicket second = rendezvous.openAsk(W); + assertTrue(second.fresh()); + assertEquals(Rendezvous.Owner.of("term_other"), rendezvous.askOwner(second.turnId()), + "a freshly-opened ask always copies whatever owner is open right now, not a stale one"); + } + + @Test + void askOwnerIsNullForAnUnknownOrLapsedTurn() { + assertNull(rendezvous.askOwner("no-such#1")); + Rendezvous.AskTicket t = rendezvous.openAsk(W); + rendezvous.closeAsk(t.turnId()); + assertNull(rendezvous.askOwner(t.turnId()), "a closed ask no longer reports an owner"); + } + + @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"); + 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), + "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 ecd83207..43fd069e 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java @@ -343,7 +343,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)); + () -> messages.answer(turnId, "config.yaml", 5000, "term_a")); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); deadline = System.currentTimeMillis() + 3000; while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) { @@ -358,6 +358,154 @@ class FleetAppAuthTest { } } + /** + * {@code POST /sessions/{id}/message} with a {@code turnId} refuses a caller whose terminal + * did not open the turn, even when that caller otherwise holds ANSWER rights, and leaves the + * turn open for the real owner to resolve. Covers the REST adapter's ANSWER gate for an + * async-send ({@code wait:false}) delegation. + */ + @Test + void restAnswerIsRefusedForADifferentCallerOnAnAsyncSendDelegationButTheRealOwnerSucceeds() throws Exception { + FakeHerdr herdr = new FakeHerdr(); + AgentControl agents = new AgentControl(herdr); + Injector injector = new Injector(agents); + Rendezvous rendezvous = new Rendezvous(); + MessageService messages = new MessageService(agents, injector, rendezvous); + + Javalin ownerApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID, Map.of("term_a", "lead-owner")); + Javalin attackerApp = startOnSharedService(messages, herdr, 9001L, Map.of("term_shell", "lead-attacker")); + try { + ObjectMapper mapper = new ObjectMapper(); + HttpResponse created = send(ownerApp.port(), "POST", "/sessions/term_target/message", + "{\"content\":\"long task\",\"wait\":false}", null); + assertEquals(202, created.statusCode(), created.body()); + String ticket = mapper.readTree(created.body()).path("ticket").asText(null); + assertNotNull(ticket, "the accepted response carried no ticket: " + created.body()); + + long deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting("term_target"), "the async send should have opened its rendezvous waiter"); + + CompletableFuture ask = CompletableFuture.supplyAsync( + () -> messages.ask("term_target", "which config file?", 5000)); + + MessageService.TaskView asking; + deadline = System.currentTimeMillis() + 3000; + do { + asking = messages.poll(ticket, null); + Thread.sleep(5); + } while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline); + assertEquals(MessageService.Phase.ASKING, asking.phase()); + String turnId = asking.turnId(); + + HttpResponse hijacked = send(attackerApp.port(), "POST", "/sessions/term_target/message", + "{\"content\":\"hijack\",\"turnId\":\"" + turnId + "\"}", null); + assertEquals(403, hijacked.statusCode(), hijacked.body()); + assertTrue(hijacked.body().contains("not_turn_owner"), + "a different caller's answer must be refused as not_turn_owner: " + hijacked.body()); + assertEquals("term_target", rendezvous.askSession(turnId), + "a refused answer must leave the ask turn open"); + + CompletableFuture> owned = CompletableFuture.supplyAsync(() -> { + try { + return send(ownerApp.port(), "POST", "/sessions/term_target/message", + "{\"content\":\"config.yaml\",\"turnId\":\"" + turnId + "\"}", null); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.resolve("term_target", "done")); + HttpResponse ownedResponse = owned.get(5, TimeUnit.SECONDS); + assertEquals(200, ownedResponse.statusCode(), ownedResponse.body()); + assertEquals("done", mapper.readTree(ownedResponse.body()).path("reply").asText(null), + "the real owner's answer must resolve the worker's turn"); + } finally { + ownerApp.stop(); + attackerApp.stop(); + } + } + + /** + * As above, but the delegation is a blocking send ({@code wait:true}) instead of a ticket — + * the owning caller's own HTTP request is the one that surfaces the worker's question and + * later carries the real answer. Covers the REST adapter's ANSWER gate for a blocking-send + * delegation. + */ + @Test + void restAnswerIsRefusedForADifferentCallerOnABlockingSendDelegationButTheRealOwnerSucceeds() throws Exception { + FakeHerdr herdr = new FakeHerdr(); + AgentControl agents = new AgentControl(herdr); + Injector injector = new Injector(agents); + Rendezvous rendezvous = new Rendezvous(); + MessageService messages = new MessageService(agents, injector, rendezvous); + + Javalin ownerApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID, Map.of("term_a", "lead-owner")); + Javalin attackerApp = startOnSharedService(messages, herdr, 9001L, Map.of("term_shell", "lead-attacker")); + try { + ObjectMapper mapper = new ObjectMapper(); + CompletableFuture> blocking = CompletableFuture.supplyAsync(() -> { + try { + return send(ownerApp.port(), "POST", "/sessions/term_target/message", + "{\"content\":\"long task\",\"wait\":true,\"timeoutMs\":5000}", null); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + + long deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting("term_target"), "the blocking send should have opened its rendezvous waiter"); + + CompletableFuture ask = CompletableFuture.supplyAsync( + () -> messages.ask("term_target", "which config file?", 5000)); + + HttpResponse questionResponse = blocking.get(5, TimeUnit.SECONDS); + assertEquals(202, questionResponse.statusCode(), questionResponse.body()); + JsonNode question = mapper.readTree(questionResponse.body()); + assertEquals("question", question.get("status").asText()); + String turnId = question.get("turnId").asText(); + + HttpResponse hijacked = send(attackerApp.port(), "POST", "/sessions/term_target/message", + "{\"content\":\"hijack\",\"turnId\":\"" + turnId + "\"}", null); + assertEquals(403, hijacked.statusCode(), hijacked.body()); + assertTrue(hijacked.body().contains("not_turn_owner"), + "a different caller's answer must be refused as not_turn_owner: " + hijacked.body()); + assertEquals("term_target", rendezvous.askSession(turnId), + "a refused answer must leave the ask turn open"); + + CompletableFuture> owned = CompletableFuture.supplyAsync(() -> { + try { + return send(ownerApp.port(), "POST", "/sessions/term_target/message", + "{\"content\":\"config.yaml\",\"turnId\":\"" + turnId + "\"}", null); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.resolve("term_target", "done")); + HttpResponse ownedResponse = owned.get(5, TimeUnit.SECONDS); + assertEquals(200, ownedResponse.statusCode(), ownedResponse.body()); + assertEquals("done", mapper.readTree(ownedResponse.body()).path("reply").asText(null), + "the real owner's answer must resolve the worker's turn"); + } finally { + ownerApp.stop(); + attackerApp.stop(); + } + } + /** * As {@link #start}, but shares {@code messages} and {@code herdr} across several app * instances bound to different pids, each returned as its own started {@link Javalin} rather From 2374de28e45d3d5068ba000033097629613771a8 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 4 Oct 2026 09:56:05 +0200 Subject: [PATCH 2/2] fleetd #715: pin that no production caller uses Rendezvous.open(String) Add a source-scrape test mirroring #718's MessageServicePollUsageTest shape, for the same residual: a convenience overload that defaults the turn's owner to null, left in place because deleting it would break 81 test-only call sites across 7 unrelated files. The scanner is exercised against a file known to hold many real one-argument rendezvous.open( calls before it is ever pointed at production, using the identical matching logic for both. A pattern that cannot find the known calls would also find none in production, and that is exactly the failure mode a -based git grep regex hit earlier on this ticket: git grep's -E engine does not treat \b as a word boundary, so that pattern silently matched nothing anywhere, in clean code and in the 81 real calls alike. --- .../fleet/msg/RendezvousOpenUsageTest.java | 187 ++++++++++++++++++ 1 file changed, 187 insertions(+) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousOpenUsageTest.java diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousOpenUsageTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousOpenUsageTest.java new file mode 100644 index 00000000..2258c016 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/RendezvousOpenUsageTest.java @@ -0,0 +1,187 @@ +package dev.ltms.fleet.msg; + +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Stream; + +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Pins that no file under {@code src/main/java} calls the fail-open + * {@link Rendezvous#open(String)} overload. That overload opens a forward waiter with no + * recorded {@link Rendezvous.Owner}, so the turn it opens can never be answered by anyone, + * not even the unnamed primary. Every production caller must go through + * {@link Rendezvous#open(String, Rendezvous.Owner)} and record an explicit owner. + * + *

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. + * + *

The scan below finds a violation by its receiver, {@code rendezvous.open(}, classified by + * argument count. The one test method here points that exact scanner at a file known to hold + * many real one-argument calls before it ever looks at production, so a scanner that stops + * matching fails loudly on the known-positive case instead of leaving a clean production result + * looking like evidence it never produced. + */ +class RendezvousOpenUsageTest { + + private static final Path PRODUCTION_SOURCE = Path.of("src/main/java"); + private static final Path KNOWN_TEST_CALLER = + Path.of("src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java"); + + /** + * A pattern that cannot find a known one-argument {@code rendezvous.open(} call would also + * find none in production -- not because production is clean, but because the pattern does + * not match the text it is supposed to catch. That failure mode is exactly what let a + * {@code \b}-based regex read "no callers" under {@code git grep -E} when 81 real ones + * existed: {@code git grep} does not treat {@code \b} as a word boundary, so the pattern + * silently matched nothing anywhere, clean code and real calls alike. This test runs the + * known-positive check first, with the same scanning method the production check then + * depends on, so that mistake fails loudly here instead of reading as a clean result. + */ + @Test + void noProductionFileCallsTheSingleArgumentOpenOverload() throws IOException { + List knownCalls = new ArrayList<>(); + int knownFilesScanned = scanForOneArgOpenCalls(KNOWN_TEST_CALLER, knownCalls); + + // CONTROL: the scan actually walked files -- a wrong root would otherwise report "found + // nothing" having looked at nothing. + assertTrue(knownFilesScanned > 0, "control failed: the scan under " + KNOWN_TEST_CALLER + + " visited zero .java files -- the path is wrong, so neither result below proves " + + "anything"); + + // CONTROL: the scanner actually finds real one-argument rendezvous.open( calls when + // pointed at a file known to hold many. If this is not satisfied, the matching logic + // itself is broken, and the production result below is the scanner failing silently, + // not production code actually being clean. + assertTrue(knownCalls.size() >= 60, "control failed: the scanner found only " + + knownCalls.size() + " one-argument rendezvous.open( call(s) in " + KNOWN_TEST_CALLER + + ", which is known to hold many -- the matching logic itself is broken: " + knownCalls); + + List violations = new ArrayList<>(); + int filesScanned = scanForOneArgOpenCalls(PRODUCTION_SOURCE, violations); + + // CONTROL: the scan actually walked files -- a wrong root would otherwise report "no + // violations found" having looked at nothing. + assertTrue(filesScanned > 0, "control failed: the scan under " + PRODUCTION_SOURCE + + " visited zero .java files -- the path is wrong, so the absence of violations " + + "below proves nothing"); + + assertTrue(violations.isEmpty(), "found a call to the fail-open Rendezvous.open(String) " + + "overload, which opens a forward waiter with no recorded owner -- record an " + + "explicit Rendezvous.Owner through open(String, Owner) instead: " + violations); + } + + private static int scanForOneArgOpenCalls(Path root, List sites) throws IOException { + List files = javaFiles(root); + for (Path file : files) { + scanFileForOpenCalls(file, sites); + } + return files.size(); + } + + private static List javaFiles(Path root) throws IOException { + try (Stream paths = Files.walk(root)) { + return paths.filter(p -> p.toString().endsWith(".java")).toList(); + } + } + + private static void scanFileForOpenCalls(Path file, List sites) throws IOException { + String source = Files.readString(file); + String needle = "rendezvous.open("; + int from = 0; + int idx; + while ((idx = source.indexOf(needle, from)) >= 0) { + int argsStart = idx + needle.length(); + String args = extractBalancedArgs(source, argsStart, file, idx); + int closeParenIndex = argsStart + args.length(); + from = closeParenIndex + 1; + + if (args.isBlank()) { + continue; // Rendezvous has no zero-argument open() -- a javadoc "rendezvous.open()" mention, not a call + } + if (topLevelCommaCount(args) == 0) { + sites.add(file + ":" + lineOf(source, idx) + " -- rendezvous.open(" + args.trim() + ")"); + } + } + } + + /** + * The text between {@code rendezvous.open(} and its matching close paren: balanced over + * nested calls, and never split by a paren or comma sitting inside a string or char literal. + */ + private static String extractBalancedArgs(String source, int start, Path file, int callIndex) { + int depth = 1; + boolean inString = false; + boolean inChar = false; + int i = start; + while (i < source.length()) { + char c = source.charAt(i); + if (inString) { + if (c == '\\') { i += 2; continue; } + if (c == '"') inString = false; + } else if (inChar) { + if (c == '\\') { i += 2; continue; } + if (c == '\'') inChar = false; + } else if (c == '"') { + inString = true; + } else if (c == '\'') { + inChar = true; + } else if (c == '(') { + depth++; + } else if (c == ')') { + depth--; + if (depth == 0) return source.substring(start, i); + } + i++; + } + throw new IllegalStateException( + "unbalanced parentheses scanning " + file + ":" + lineOf(source, callIndex)); + } + + /** + * Commas at paren/bracket/brace depth zero, skipping string and char literals -- the argument + * separators a human reader would see, not every comma character in the text. + */ + private static int topLevelCommaCount(String args) { + int depth = 0; + int commas = 0; + boolean inString = false; + boolean inChar = false; + int i = 0; + while (i < args.length()) { + char c = args.charAt(i); + if (inString) { + if (c == '\\') { i += 2; continue; } + if (c == '"') inString = false; + } else if (inChar) { + if (c == '\\') { i += 2; continue; } + if (c == '\'') inChar = false; + } else if (c == '"') { + inString = true; + } else if (c == '\'') { + inChar = true; + } else if (c == '(' || c == '[' || c == '{') { + depth++; + } else if (c == ')' || c == ']' || c == '}') { + depth--; + } else if (c == ',' && depth == 0) { + commas++; + } + i++; + } + return commas; + } + + private static int lineOf(String source, int index) { + int line = 1; + for (int i = 0; i < index; i++) { + if (source.charAt(i) == '\n') line++; + } + return line; + } +}