From 0ba597e394cd44cd49697c72b18a8186cea634ce Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 4 Oct 2026 07:23:29 +0200 Subject: [PATCH 1/2] fleetd #705: close the REST ticket-poll door and pin both handlers' caller-terminal wiring GET /tasks/{ticket} now resolves the caller the same way allow(...) does and threads that terminal into MessageService.poll(ticket, callerTerminal) instead of the no-check overload, so a worker can no longer read a ticket a different session created over REST. Adds a source-scrape guard (with its own control assertion) for both the fleet_poll MCP handler and this REST route, plus a behavioural test driving GET /tasks/{ticket} with three differently-resolved callers against one shared MessageService. --- .../java/dev/ltms/fleet/rest/FleetApp.java | 3 +- .../dev/ltms/fleet/mcp/FleetMcpAuthzTest.java | 27 +++++ .../dev/ltms/fleet/rest/FleetAppAuthTest.java | 100 ++++++++++++++++++ 3 files changed, 129 insertions(+), 1 deletion(-) 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 824f5ff..c9c3c7e 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -894,7 +894,8 @@ public final class FleetApp { if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) { return; } - MessageService.TaskView v = messages.poll(ctx.pathParam("ticket")); + Principal caller = ctx.attribute(CALLER); + MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.terminal()); 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/mcp/FleetMcpAuthzTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java index cb6db34..5f154a2 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java @@ -418,6 +418,33 @@ class FleetMcpAuthzTest { + "calling, not pass a literal boolean -- block: " + handlerBlock); } + /** + * {@code fleet_poll{ticket}} must thread the calling connection's own terminal into + * {@link MessageService#poll(String, String)}, so a worker cannot read a ticket a different + * session created. + */ + @Test + void theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll() throws Exception { + String source = Files.readString(MCP_SOURCE); + + int start = source.indexOf("pollHandler ="); + assertTrue(start >= 0, "could not find the fleet_poll handler (pollHandler) in " + MCP_SOURCE + + " -- the scrape has stopped matching, fix the anchor before trusting this test"); + int end = source.indexOf("ackHandler =", start); + assertTrue(end > start, "could not find the handler declared after pollHandler to bound the scrape"); + String handlerBlock = source.substring(start, end); + + // CONTROL: the block we scraped really does call poll(...) -- if this fails, the anchors + // above moved and the assertions below would otherwise pass on nothing. + assertTrue(handlerBlock.contains("poll(messages,"), + "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 " + + "it or pass a literal null -- block: " + handlerBlock); + } + // --- which action each tool hands the gate (fleetd #272) ------------------------------------ /** 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 735c116..3cadf92 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java @@ -139,6 +139,106 @@ class FleetAppAuthTest { assertThrows(IllegalArgumentException.class, () -> FleetApp.routeAction("GET /healthz")); } + // --- GET /tasks/{ticket} must not be the no-check overload ---------------------------------- + + /** + * {@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 + * no-check overload that ignores who is asking. + */ + @Test + void theTaskStatusRouteActuallyThreadsTheCallersTerminalIntoPoll() throws Exception { + String source = Files.readString(REST_SOURCE); + + int start = source.indexOf("private void taskStatus(Context ctx) {"); + assertTrue(start >= 0, "could not find taskStatus in " + REST_SOURCE + + " -- the scrape has stopped matching, fix the anchor before trusting this test"); + int end = source.indexOf("private static void herdrError(Context ctx, HerdrException e) {", start); + assertTrue(end > start, "could not find the method declared after taskStatus to bound the scrape"); + String handlerBlock = source.substring(start, end); + + // CONTROL: the block we scraped really does call messages.poll(...) -- if this fails, the + // anchors above moved and the assertions below would otherwise pass on nothing. + assertTrue(handlerBlock.contains("messages.poll("), + "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(...), " + + "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 " + + "second, separate resolution path -- block: " + handlerBlock); + } + + /** + * {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket, + * while the creating worker and the unnamed primary both still read it. The ticket is minted + * directly on the shared {@link MessageService}, the same way {@code MessageServiceTest} + * drives {@link MessageService#poll(String, String)}, so this exercises only the REST poll + * route's own handling of the ownership already recorded on the ticket. + */ + @Test + void restPollRefusesADifferentWorkerButAllowsTheCreatorAndTheUnnamedPrimary() throws Exception { + FakeHerdr herdr = new FakeHerdr(); + AgentControl agents = new AgentControl(herdr); + Injector injector = new Injector(agents); + MessageService messages = new MessageService(agents, injector, new Rendezvous()); + + Javalin creatorApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID); // -> term_a + 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"); + + HttpResponse refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null); + assertEquals(200, refused.statusCode()); + assertTrue(refused.body().contains("forbidden"), + "a different worker's terminal must be refused, not shown the ticket: " + refused.body()); + assertFalse(refused.body().contains("\"reply\""), + "a refusal must never carry reply text: " + refused.body()); + + HttpResponse own = send(creatorApp.port(), "GET", "/tasks/" + ticket, null, null); + assertEquals(200, own.statusCode()); + assertFalse(own.body().contains("forbidden"), + "the creating worker must read its own ticket: " + own.body()); + + HttpResponse primary = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null); + assertEquals(200, primary.statusCode()); + assertFalse(primary.body().contains("forbidden"), + "the unnamed primary must read any ticket: " + primary.body()); + } finally { + creatorApp.stop(); + otherWorkerApp.stop(); + primaryApp.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 + * than through the shared {@code app} field, so several differently-resolved callers can + * poll the same ticket. + */ + private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid) { + FleetConfig.Profile wcfg = new FleetConfig.Profile( + "ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null, + "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null); + AgentControl agents = new AgentControl(herdr); + ClaudeCodeLauncher workers = new ClaudeCodeLauncher( + agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), + Map.of(wcfg.profile(), wcfg), wcfg.profile(), + k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null); + SessionManager sessions = new SessionManager(workers, new FakeWorktrees()); + ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid); + CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, + Map::of, new MemberRegistry(null)); + Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox()); + + return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null, + callers, appMetrics).build().start("127.0.0.1", 0); + } + /** * fleetd #669 Unit A: {@code POST /sessions/{id}/message} is two call shapes behind one route, * mirroring {@code fleet_send}'s MCP-side split into {@link Authz.Action#SEND} and {@link From 29a2f97c25a7af1c0cfecedd76ec8135d3f97ba7 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 4 Oct 2026 07:41:12 +0200 Subject: [PATCH 2/2] fleetd: thread the creating caller's terminal into the REST fire-and-poll send path sendMessage's wait:false branch now records the resolved caller's own terminal as the ticket's creatorTerminal, the same way taskStatus already resolves its caller, so a REST-created ticket's own creator can still poll it under the ownership check that now gates GET /tasks/{ticket}. --- .../java/dev/ltms/fleet/rest/FleetApp.java | 3 +- .../dev/ltms/fleet/rest/FleetAppAuthTest.java | 50 ++++++++++++++++++- 2 files changed, 51 insertions(+), 2 deletions(-) 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 c9c3c7e..c7ee1f7 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -694,7 +694,8 @@ public final class FleetApp { if (!wait) { // Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}. - String ticket = messages.sendAsync(id, content); + Principal caller = ctx.attribute(CALLER); + String ticket = messages.sendAsync(id, content, null, caller == null ? null : caller.terminal()); ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted")); return; } 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 3cadf92..2b93b4e 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java @@ -214,6 +214,44 @@ class FleetAppAuthTest { } } + /** + * {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating + * caller's own terminal on the ticket it returns, so that caller can still poll its own + * ticket over REST, while a different terminal is refused. + */ + @Test + void restSendAsyncRecordsTheCreatingCallersTerminalSoItCanStillPollItsOwnTicket() throws Exception { + FakeHerdr herdr = new FakeHerdr(); + AgentControl agents = new AgentControl(herdr); + Injector injector = new Injector(agents); + MessageService messages = new MessageService(agents, injector, new Rendezvous()); + + Javalin leadApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID, Map.of("term_a", "lead-x")); + Javalin otherWorkerApp = startOnSharedService(messages, herdr, 9001L); // -> term_shell + try { + ObjectMapper mapper = new ObjectMapper(); + HttpResponse created = send(leadApp.port(), "POST", "/sessions/term_a/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()); + + HttpResponse own = send(leadApp.port(), "GET", "/tasks/" + ticket, null, null); + assertEquals(200, own.statusCode()); + assertFalse(own.body().contains("forbidden"), + "the session that created the ticket over REST must be able to poll it: " + own.body()); + + HttpResponse refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null); + assertEquals(200, refused.statusCode()); + assertTrue(refused.body().contains("forbidden: this ticket was created by a different session"), + "a different terminal must still be refused with the ownership detail, not some " + + "other rejection: " + refused.body()); + } finally { + leadApp.stop(); + otherWorkerApp.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 @@ -221,6 +259,16 @@ class FleetAppAuthTest { * poll the same ticket. */ private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid) { + return startOnSharedService(messages, herdr, pid, Map.of()); + } + + /** + * As {@link #startOnSharedService(MessageService, FakeHerdr, long)}, but {@code leadTerminals} + * resolves the given pid's terminal to a named lead (a caller with SEND permission) instead of + * a plain worker, for a test that needs a terminal-bearing caller able to create a ticket. + */ + private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid, + Map leadTerminals) { FleetConfig.Profile wcfg = new FleetConfig.Profile( "ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null, "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null); @@ -232,7 +280,7 @@ class FleetAppAuthTest { SessionManager sessions = new SessionManager(workers, new FakeWorktrees()); ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid); CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, - Map::of, new MemberRegistry(null)); + () -> leadTerminals, new MemberRegistry(null)); Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox()); return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,