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 2507a058..2082ec87 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -514,7 +514,7 @@ public final class FleetMcp { (exchange, req) -> { McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null); if (denied != null) return denied; - return status(messages, str(req.arguments(), "sessionId")); + return status(messages, str(req.arguments(), "sessionId"), callerTerminal(exchange)); }; BiFunction pollHandler = (exchange, req) -> { @@ -1335,16 +1335,20 @@ public final class FleetMcp { /** * {@code fleet_status}: the live lifecycle status of a worker session, plus — when the worker - * is paused mid-turn in an async {@code fleet_ask} (CB-582) — the open question and how to - * answer it, so a lead on its normal poll cadence does not need the ticket to notice. + * is paused mid-turn in an async {@code fleet_ask} — the open question and how to answer it, so + * a lead on its normal poll cadence does not need the ticket to notice. The question, its + * {@code turnId} and its ticket id are shown only to the caller whose terminal created that + * delegation, or to a caller with no terminal at all (the unnamed primary); any other caller + * still sees the base status. {@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 status(MessageService messages, String sessionId) { + static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerTerminal) { if (isBlank(sessionId)) { return error("sessionId is required"); } try { String base = messages.status(sessionId).name().toLowerCase(); - MessageService.PendingAsk ask = messages.pendingAsk(sessionId); + MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerTerminal); if (ask == null) { return text(base); } 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 31db43a7..da466c93 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -1754,16 +1754,17 @@ public final class MessageService { } /** - * The question {@code workerSession} is currently paused on via {@code fleet_ask}, if any - * (CB-582) — {@code fleet_status} uses this to show a pending question without the caller - * needing the ticket. {@code null} when the session has no open async question (including a - * session mid a blocking {@code fleet_ask}, which has no {@link Task} to look up — see - * {@link PendingAsk}). + * The question {@code workerSession} is currently paused on via {@code fleet_ask}, if any — + * {@code fleet_status} uses this to show a pending question without the caller needing the + * ticket. {@code null} when the session has no open async question (including a session mid a + * blocking {@code fleet_ask}, which has no {@link Task} to look up — see + * {@link PendingAsk}), or when {@code callerTerminal} does not own the task the question + * belongs to (see {@link #ownsTicket(Task, String)}). */ - public PendingAsk pendingAsk(String workerSession) { + public PendingAsk pendingAsk(String workerSession, String callerTerminal) { for (Task task : tasks.values()) { Reply q = task.question; - if (q != null && workerSession.equals(task.target)) { + if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerTerminal)) { return new PendingAsk(task.ticket, q.text(), q.turnId()); } } 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 c7ee1f73..93254acb 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -875,10 +875,12 @@ public final class FleetApp { body.put("sessionId", id); body.put("status", messages.status(id).name().toLowerCase()); body.put("ready", deliverable.test(id)); - // CB-582: a worker paused mid-turn in an async fleet_ask is otherwise invisible to a - // status poll — surface the open question and how to answer it, same as fleet_poll's - // Phase.ASKING view. - MessageService.PendingAsk ask = messages.pendingAsk(id); + // A worker paused mid-turn in an async fleet_ask is otherwise invisible to a status + // poll — surface the open question and how to answer it, same as fleet_poll's + // Phase.ASKING view, but only to the caller whose terminal created that delegation, or + // to a caller with no terminal at all (the unnamed primary). + Principal caller = ctx.attribute(CALLER); + MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.terminal()); if (ask != null) { body.put("question", ask.question()); body.put("turnId", ask.turnId()); 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 6ff21411..c15acd33 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java @@ -530,6 +530,33 @@ class FleetMcpAuthzTest { + "it or pass a literal null -- block: " + handlerBlock); } + /** + * {@code fleet_status} must thread the calling connection's own terminal into + * {@link FleetMcp#status(MessageService, String, String)}, so a caller that did not create a + * worker's open delegation cannot read its pending question through the status handler either. + */ + @Test + void theFleetStatusHandlerActuallyThreadsCallerTerminalIntoStatus() throws Exception { + String source = Files.readString(MCP_SOURCE); + + int start = source.indexOf("statusHandler ="); + assertTrue(start >= 0, "could not find the fleet_status handler (statusHandler) in " + MCP_SOURCE + + " -- the scrape has stopped matching, fix the anchor before trusting this test"); + int end = source.indexOf("pollHandler =", start); + assertTrue(end > start, "could not find the handler declared after statusHandler to bound the scrape"); + String handlerBlock = source.substring(start, end); + + // CONTROL: the block we scraped really does call status(...) -- if this fails, the anchors + // above moved and the assertions below would otherwise pass on nothing. + assertTrue(handlerBlock.contains("status(messages,"), + "control failed: the scraped statusHandler block contains no status(messages, ...) " + + "call at all -- the anchors have drifted, this test is not testing what it claims to"); + + assertTrue(handlerBlock.contains("callerTerminal(exchange)"), + "the fleet_status handler must thread callerTerminal(exchange) into status(...), 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/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 42e04b62..74c49f0c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -1860,14 +1860,14 @@ class FleetMcpTest { FakeHerdr blocked = new FakeHerdr().agentStatus("blocked"); AgentControl blockedAgents = new AgentControl(blocked); McpSchema.CallToolResult res = FleetMcp.status( - new MessageService(blockedAgents, new Injector(blockedAgents), rendezvous), "term_a"); + new MessageService(blockedAgents, new Injector(blockedAgents), rendezvous), "term_a", null); assertNotEquals(Boolean.TRUE, res.isError()); assertEquals("blocked", textOf(res)); } /** - * CB-582: a lead polling {@code fleet_status} on its normal cadence — not {@code fleet_poll} - * — must also see a worker's open async {@code fleet_ask} question, since the reverse-rendezvous + * A lead polling {@code fleet_status} on its normal cadence — not {@code fleet_poll} — must + * also see a worker's open async {@code fleet_ask} question, since the reverse-rendezvous * window it opened with is far shorter than that cadence. */ @Test @@ -1890,7 +1890,7 @@ class FleetMcpTest { } while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline); assertEquals(MessageService.Phase.ASKING, asking.phase()); - McpSchema.CallToolResult res = FleetMcp.status(messages, T); + McpSchema.CallToolResult res = FleetMcp.status(messages, T, null); assertNotEquals(Boolean.TRUE, res.isError()); String out = textOf(res); assertTrue(out.startsWith("idle"), "the live status must still lead the text: " + out); @@ -1911,6 +1911,63 @@ class FleetMcpTest { answer.get(5, TimeUnit.SECONDS); } + /** + * {@code fleet_status}'s pending-ask block (the question, its {@code turnId} and its ticket) + * is shown only to the caller whose terminal created the delegation, or to a caller with no + * terminal at all (the unnamed primary) — a different terminal-bearing caller still sees the + * base status line, but none of the pending-ask fields. + */ + @Test + void statusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception { + String ticket = messages.sendAsync(T, "task that asks", null, "term_creator"); + long deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting(T), "sendAsync should have opened its rendezvous waiter"); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + + MessageService.TaskView asking; + deadline = System.currentTimeMillis() + 3000; + do { + asking = messages.poll(ticket); + Thread.sleep(5); + } while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline); + assertEquals(MessageService.Phase.ASKING, asking.phase()); + + String other = textOf(FleetMcp.status(messages, T, "term_other")); + assertTrue(other.startsWith("idle"), "the base status must still be shown: " + other); + assertFalse(other.contains("which config file?"), + "a non-creating caller must not see the question text: " + other); + assertFalse(other.contains(asking.turnId()), + "a non-creating caller must not see the turnId: " + other); + assertFalse(other.contains(ticket), + "a non-creating caller must not see the ticket: " + other); + + String creator = textOf(FleetMcp.status(messages, T, "term_creator")); + assertTrue(creator.contains("which config file?"), "the creator must see the question: " + creator); + assertTrue(creator.contains(asking.turnId()), "the creator must see the turnId: " + creator); + assertTrue(creator.contains(ticket), "the creator must see the ticket: " + creator); + + String unnamed = textOf(FleetMcp.status(messages, T, null)); + assertTrue(unnamed.contains("which config file?"), + "a caller with no terminal (the unnamed primary) must see the question: " + unnamed); + + // 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)); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.resolve(T, "done")); + answer.get(5, TimeUnit.SECONDS); + } + // --- fleet_whoami: the caller's own identity, so an agent never has to guess its role ------- @Test 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 9b2ad571..912da690 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -1875,20 +1875,20 @@ class MessageServiceTest { } } - // --- CB-582: fleet_status pendingAsk() ------------------------------------------------------ + // --- fleet_status pendingAsk() ----------------------------------------------------------- @Test void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception { - assertNull(messages.pendingAsk(T), "no async ticket at all -> no pending ask"); + assertNull(messages.pendingAsk(T, null), "no async ticket at all -> no pending ask"); String ticket = messages.sendAsync(T, "long task"); awaitWaiting(); - assertNull(messages.pendingAsk(T), "a plain pending delegation is not a question"); + assertNull(messages.pendingAsk(T, null), "a plain pending delegation is not a question"); injectDelivery(); assertTrue(rendezvous.resolve(T, "done")); awaitTicketPhase(ticket, MessageService.Phase.DONE); - assertNull(messages.pendingAsk(T), "a finished ticket carries no open question either"); + assertNull(messages.pendingAsk(T, null), "a finished ticket carries no open question either"); } @Test @@ -1901,7 +1901,7 @@ class MessageServiceTest { CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); - MessageService.PendingAsk pending = messages.pendingAsk(T); + MessageService.PendingAsk pending = messages.pendingAsk(T, null); assertNotNull(pending, "fleet_status should see the open question"); assertEquals(ticket, pending.ticket()); assertEquals("which config file?", pending.question()); @@ -1910,7 +1910,69 @@ class MessageServiceTest { CompletableFuture answer = CompletableFuture.supplyAsync( () -> messages.answer(asking.turnId(), "config.yaml", 5000)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); - assertNull(messages.pendingAsk(T), "an answered question is no longer pending"); + assertNull(messages.pendingAsk(T, null), "an answered question is no longer pending"); + awaitWaiting(); + assertTrue(rendezvous.resolve(T, "done")); + assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome()); + } + + /** + * A caller's own terminal must match the terminal that created the delegation to see its + * pending question; a different terminal-bearing caller sees nothing, and a caller with no + * terminal at all (the unnamed primary) always sees it. + */ + @Test + void pendingAskGatesTheQuestionByTheDelegationsCreatorTerminal() throws Exception { + String ticket = messages.sendAsync(T, "task that asks", null, "term_creator"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); + + assertNull(messages.pendingAsk(T, "term_other"), + "a caller whose terminal did not create the delegation must not see the question"); + + MessageService.PendingAsk own = messages.pendingAsk(T, "term_creator"); + assertNotNull(own, "the creating caller must see its own open question"); + assertEquals("which config file?", own.question()); + + MessageService.PendingAsk unnamed = messages.pendingAsk(T, null); + assertNotNull(unnamed, "a caller with no terminal (the unnamed primary) must always see the question"); + assertEquals("which config file?", unnamed.question()); + + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> messages.answer(asking.turnId(), "config.yaml", 5000)); + 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()); + } + + /** + * A task created with no recorded creator terminal (a short {@code sendAsync} overload) must + * not hand its open question to any caller that does have a terminal — only a caller with no + * terminal at all may still see it. + */ + @Test + void pendingAskDeniesATerminalBearingCallerWhenTheTaskRecordsNoCreator() throws Exception { + String ticket = messages.sendAsync(T, "task that asks"); // no creatorTerminal recorded + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); + + assertNull(messages.pendingAsk(T, "term_someone"), + "a terminal-bearing caller must not see a question whose task records no creator"); + assertNotNull(messages.pendingAsk(T, null), + "the unnamed primary must still see it even with no recorded creator"); + + CompletableFuture answer = CompletableFuture.supplyAsync( + () -> messages.answer(asking.turnId(), "config.yaml", 5000)); + 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()); 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 2b93b4ea..ecd83207 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppAuthTest.java @@ -2,6 +2,7 @@ package dev.ltms.fleet.rest; import ch.qos.logback.classic.Level; import ch.qos.logback.classic.spi.ILoggingEvent; +import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import dev.ltms.fleet.auth.CallerResolver; import dev.ltms.fleet.auth.Authz; @@ -39,6 +40,8 @@ import java.util.List; import java.util.Locale; import java.util.Map; import java.util.Set; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; import java.util.function.Predicate; import java.util.regex.Matcher; import java.util.regex.Pattern; @@ -171,6 +174,36 @@ class FleetAppAuthTest { + "second, separate resolution path -- block: " + handlerBlock); } + /** + * {@code GET /sessions/{id}/status} must resolve its caller the same way {@code allow(...)} + * does and thread that terminal into {@link MessageService#pendingAsk(String, String)}, not + * the no-check overload that ignores who is asking. + */ + @Test + void theSessionStatusRouteActuallyThreadsTheCallersTerminalIntoPendingAsk() throws Exception { + String source = Files.readString(REST_SOURCE); + + int start = source.indexOf("private void sessionStatus(Context ctx) {"); + assertTrue(start >= 0, "could not find sessionStatus in " + REST_SOURCE + + " -- the scrape has stopped matching, fix the anchor before trusting this test"); + int end = source.indexOf("private void taskStatus(Context ctx) {", start); + assertTrue(end > start, "could not find the method declared after sessionStatus to bound the scrape"); + String handlerBlock = source.substring(start, end); + + // CONTROL: the block we scraped really does call messages.pendingAsk(...) -- if this fails, + // the anchors above moved and the assertions below would otherwise pass on nothing. + assertTrue(handlerBlock.contains("messages.pendingAsk("), + "control failed: the scraped sessionStatus block contains no messages.pendingAsk( " + + "call at all -- the anchors have drifted, this test is not testing what it claims to"); + + assertTrue(handlerBlock.contains("caller.terminal()"), + "the sessionStatus route must thread the resolved caller's terminal into " + + "messages.pendingAsk(...), not the no-check overload -- block: " + handlerBlock); + assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"), + "the sessionStatus 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 @@ -252,6 +285,79 @@ class FleetAppAuthTest { } } + /** + * {@code GET /sessions/{id}/status} shows a worker's pending {@code fleet_ask} question, its + * {@code turnId} and its ticket only to the caller whose terminal created that delegation, or + * to a caller with no terminal at all (the unnamed primary) — a different terminal-bearing + * caller still sees the base status line, but none of the pending-ask fields. + */ + @Test + void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() 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 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 { + ObjectMapper mapper = new ObjectMapper(); + String ticket = messages.sendAsync("term_target", "task that asks", null, "term_a"); + long deadline = System.currentTimeMillis() + 3000; + while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting("term_target"), "sendAsync 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(); + + JsonNode other = mapper.readTree( + send(otherWorkerApp.port(), "GET", "/sessions/term_target/status", null, null).body()); + assertEquals("idle", other.get("status").asText(), "the base status must still be shown"); + assertFalse(other.has("question"), "a non-creating caller must not see the question: " + other); + assertFalse(other.has("turnId"), "a non-creating caller must not see the turnId: " + other); + assertFalse(other.has("ticket"), "a non-creating caller must not see the ticket: " + other); + + JsonNode own = mapper.readTree( + send(creatorApp.port(), "GET", "/sessions/term_target/status", null, null).body()); + assertEquals("which config file?", own.get("question").asText(), "the creator must see the question"); + assertEquals(turnId, own.get("turnId").asText(), "the creator must see the turnId"); + assertEquals(ticket, own.get("ticket").asText(), "the creator must see the ticket"); + + JsonNode primary = mapper.readTree( + send(primaryApp.port(), "GET", "/sessions/term_target/status", null, null).body()); + assertEquals("which config file?", primary.get("question").asText(), + "a caller with no terminal (the unnamed primary) must see the question"); + + // 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)); + 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")); + answer.get(5, TimeUnit.SECONDS); + } 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