Merge PR #716: fleetd #705 — gate the REST ticket routes by the creating caller
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 1m57s

Three parts. GET /tasks/{ticket} passed no caller, so it used the overload that skips the ownership
check and any worker could read any ticket; it now passes the caller resolved from the same CALLER
attribute the authorization gate reads. The wait:false send path recorded no creator terminal, so a
REST-created ticket matched no terminal-bearing caller and its own creator was refused; it now
records one. Both handlers gained a scrape guard with its own control assertion.

Resolved one conflict in FleetMcpAuthzTest by keeping both sides: PR #717 and this branch each
appended tests at the same point. The test count is the check on that resolution — 2014 + 3 + 1 =
2018, so no test was dropped.

Verified in a throwaway worktree off main: Tests run: 2018, Failures: 0, BUILD SUCCESS. Three
mutations, each confirmed live with mvn -o compile before the suite ran. FleetMcp:527 is the one
that survived before this work and now kills theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll.
FleetApp:698 kills the new creator test. FleetApp:899 kills three, including the behavioural test.
Every file restored byte-identical.
This commit is contained in:
Dai Ha
2026-10-04 07:43:23 +02:00
3 changed files with 179 additions and 2 deletions
@@ -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;
}
@@ -894,7 +895,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;
@@ -503,6 +503,33 @@ class FleetMcpAuthzTest {
+ "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) ------------------------------------
/**
@@ -139,6 +139,154 @@ 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<String> 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<String> 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<String> primary = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, primary.statusCode());
assertFalse(primary.body().contains("forbidden"),
"the unnamed primary must read any ticket: " + primary.body());
} finally {
creatorApp.stop();
otherWorkerApp.stop();
primaryApp.stop();
}
}
/**
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
* caller's own terminal on the ticket it returns, so that caller can still poll its own
* 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<String> 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<String> 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<String> 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
* 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) {
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<String, String> 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);
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,
() -> 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,
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