|
|
|
@@ -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
|
|
|
|
|