Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 29a2f97c25 | |||
| 0ba597e394 | |||
| c468953963 | |||
| 25d53e6ef7 | |||
| a9a37af957 | |||
| b205bcc2aa | |||
| 11eccc3a1b |
@@ -249,17 +249,6 @@ public final class CallerResolver {
|
||||
return leadTerminals.get();
|
||||
}
|
||||
|
||||
/**
|
||||
* The currently-recognised architect slots, {@code terminal_id → slot name} (CB-548).
|
||||
*
|
||||
* <p>Read from the same supplier {@link #resolve} consults, so a slot that is <em>listed</em>
|
||||
* here but would not <em>resolve</em> (or the reverse) cannot drift apart. Live for the same
|
||||
* reason as {@link #leads()}.
|
||||
*/
|
||||
public Map<String, String> members() {
|
||||
return architectTerminals.get();
|
||||
}
|
||||
|
||||
/**
|
||||
* The currently-recognised collaborator tabs, {@code terminal_id → name}.
|
||||
*
|
||||
|
||||
@@ -490,7 +490,7 @@ public final class FleetMcp {
|
||||
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
|
||||
// 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())
|
||||
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
|
||||
};
|
||||
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
|
||||
@@ -524,7 +524,7 @@ public final class FleetMcp {
|
||||
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
|
||||
if (denied != null) return denied;
|
||||
return poll(messages, leadChannel, str(a, "ticket"), target, coordId);
|
||||
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
|
||||
};
|
||||
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
|
||||
// Acking removes a reply from the inbox, so it is a drain, not a read.
|
||||
@@ -972,6 +972,17 @@ public final class FleetMcp {
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
|
||||
Runnable onAccepted, Set<String> profiles) {
|
||||
return sendAsync(messages, sessionId, content, onAccepted, profiles, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
|
||||
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
|
||||
* result back to that same caller — see {@link MessageService#poll(String, String)}.
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
|
||||
Runnable onAccepted, Set<String> profiles,
|
||||
String creatorTerminal) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
@@ -979,7 +990,7 @@ public final class FleetMcp {
|
||||
if (targetError != null) {
|
||||
return targetError;
|
||||
}
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted);
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
|
||||
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
|
||||
}
|
||||
|
||||
@@ -1158,9 +1169,24 @@ public final class FleetMcp {
|
||||
* exclusively of {@code ticket}/{@code target} — see {@link #pollAction}'s javadoc for why this
|
||||
* is a different inbox (this daemon's own {@link LeadChannel}) that authorizes differently
|
||||
* ({@link Authz.Action#COORD_READ}, primary-only) from either of the original two branches.
|
||||
*
|
||||
* <p>Does not check who owns {@code ticket} — see the overload that takes {@code
|
||||
* callerTerminal} for that. Callers that do not resolve a caller terminal (tests, or a surface
|
||||
* with no connection-based identity) use this one.
|
||||
*/
|
||||
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
|
||||
String target, String coordId) {
|
||||
return poll(messages, leadChannel, ticket, target, coordId, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
|
||||
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
|
||||
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
|
||||
* client-supplied value.
|
||||
*/
|
||||
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
|
||||
String target, String coordId, String callerTerminal) {
|
||||
if (!isBlank(coordId)) {
|
||||
return pollHeldPeerMail(leadChannel, coordId);
|
||||
}
|
||||
@@ -1174,7 +1200,7 @@ public final class FleetMcp {
|
||||
if (isBlank(ticket)) {
|
||||
return error("ticket (or target) is required");
|
||||
}
|
||||
MessageService.TaskView v = messages.poll(ticket);
|
||||
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
|
||||
if (v == null) {
|
||||
return error("unknown ticket: " + ticket + " (never issued, or expired)");
|
||||
}
|
||||
|
||||
@@ -275,10 +275,19 @@ public final class MessageService {
|
||||
* {@link #abandon}) can never match again regardless of this flag's value.
|
||||
*/
|
||||
private volatile boolean askTimedOut;
|
||||
/**
|
||||
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
|
||||
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
|
||||
* created through an overload that does not record one. {@link #poll(String, String)}
|
||||
* compares a polling caller's own terminal against this field before handing back the
|
||||
* ticket's state.
|
||||
*/
|
||||
private final String creatorTerminal;
|
||||
|
||||
private Task(String ticket, String target, LongSupplier nowNanos) {
|
||||
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
|
||||
this.ticket = ticket;
|
||||
this.target = target;
|
||||
this.creatorTerminal = creatorTerminal;
|
||||
this.createdNanos = nowNanos.getAsLong();
|
||||
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
|
||||
}
|
||||
@@ -1279,8 +1288,20 @@ public final class MessageService {
|
||||
* @return the ticket to poll for the eventual result
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted) {
|
||||
return sendAsync(target, content, onAccepted, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
|
||||
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
|
||||
* differs from this one; {@code null} records no owner (a caller with no terminal — the
|
||||
* unnamed primary — is always allowed to poll the result regardless).
|
||||
*
|
||||
* @return the ticket to poll for the eventual result
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
|
||||
String ticket = "task-" + ticketSeq.incrementAndGet();
|
||||
Task task = new Task(ticket, target, nowNanos);
|
||||
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
|
||||
tasks.put(ticket, task);
|
||||
if (pushLoop != null) {
|
||||
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
||||
@@ -1343,15 +1364,31 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket;
|
||||
* otherwise a {@link Phase#PENDING} view (with the live worker status as detail), a
|
||||
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
|
||||
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
|
||||
* checked, so this overload must only be used where the caller's identity is otherwise
|
||||
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
|
||||
*/
|
||||
public TaskView poll(String ticket) {
|
||||
return poll(ticket, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
|
||||
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
|
||||
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
|
||||
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
|
||||
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
|
||||
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
|
||||
*/
|
||||
public TaskView poll(String ticket, String callerTerminal) {
|
||||
Task task = tasks.get(ticket);
|
||||
if (task == null) {
|
||||
return null;
|
||||
}
|
||||
if (!ownsTicket(task, callerTerminal)) {
|
||||
return new TaskView(ticket, Phase.FAILED, null, null,
|
||||
"forbidden: this ticket was created by a different session", null);
|
||||
}
|
||||
CompletableFuture<Reply> f = task.future;
|
||||
if (!f.isDone()) {
|
||||
Reply question = task.question;
|
||||
@@ -1387,6 +1424,17 @@ public final class MessageService {
|
||||
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
|
||||
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
|
||||
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
|
||||
* equal the terminal recorded on the task; a task with no recorded terminal matches no
|
||||
* terminal-bearing caller.
|
||||
*/
|
||||
private static boolean ownsTicket(Task task, String callerTerminal) {
|
||||
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
|
||||
}
|
||||
|
||||
/**
|
||||
* Test seam only — carries no production behaviour, and nothing in this class calls it;
|
||||
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -460,7 +460,6 @@ class CallerResolverTest {
|
||||
assertTrue(members.bind("architect:lead-designer", "term_a")); // the later lifecycle binds the slot
|
||||
|
||||
assertEquals(Role.ARCHITECT, r.resolve("127.0.0.1", 42, null).role());
|
||||
assertEquals("architect:lead-designer", r.members().get("term_a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -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) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -865,11 +865,82 @@ class MessageServiceTest {
|
||||
+ "follow-up send would come back BUSY instead of timing out on its own work");
|
||||
}
|
||||
|
||||
/**
|
||||
* Drive an async send on {@code T} to a resolved reply, polling as {@code owner} until the
|
||||
* ticket reports {@link MessageService.Phase#DONE} (or the 2s deadline runs out). {@code owner}
|
||||
* must be a terminal this ticket's creator check actually accepts, or this loops to the
|
||||
* deadline and returns a non-{@code DONE} view.
|
||||
*/
|
||||
private MessageService.TaskView driveAsyncTicketToDone(String ticket, String owner) throws InterruptedException {
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (view == null || view.phase() != MessageService.Phase.DONE) {
|
||||
if (System.currentTimeMillis() >= deadline) break;
|
||||
view = messages.poll(ticket, owner);
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
return view;
|
||||
}
|
||||
|
||||
@Test
|
||||
void pollReturnsNullForAnUnknownTicket() {
|
||||
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
|
||||
}
|
||||
|
||||
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
|
||||
|
||||
@Test
|
||||
void pollByAnotherTerminalIsRefused() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_a");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "secret async result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, "term_a");
|
||||
assertNotNull(owner, "the creator must still be able to read its own ticket");
|
||||
assertEquals(MessageService.Phase.DONE, owner.phase());
|
||||
|
||||
MessageService.TaskView refused = messages.poll(ticket, "term_b");
|
||||
assertNotNull(refused, "a different terminal gets a refusal, not silence");
|
||||
assertNotEquals(MessageService.Phase.DONE, refused.phase(),
|
||||
"a different terminal must never see the ticket as DONE");
|
||||
assertNull(refused.reply(), "a refusal must never carry the reply text");
|
||||
assertFalse(String.valueOf(refused).contains("secret async result"),
|
||||
"the reply text must not appear anywhere in the refused view");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
|
||||
|
||||
// callerTerminal == null is the unnamed primary (resolved by token or loopback trust, with
|
||||
// no herdr pane) — it must read a ticket a terminal-bearing lead created.
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
|
||||
assertNotNull(view, "the unnamed primary must be able to read any ticket");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("primary-visible result", view.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void creatorReadsItsOwnTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "own result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, "term_creator");
|
||||
assertNotNull(view, "the ticket's own creator must be able to read it");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("own result", view.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void pollReportsACompletedTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user