Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 0f5985b419 | |||
| 6e06058b07 | |||
| e5f4fb81ab | |||
| d28ab0968b | |||
| 9a64d42599 | |||
| f4e0ca41e6 | |||
| 29a2f97c25 | |||
| 95311c6e8e | |||
| 0ba597e394 |
@@ -168,7 +168,7 @@ you decide.
|
||||
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
|
||||
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
|
||||
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
|
||||
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — but **`fleet_list` does not report collaborators**, so you cannot discover one: it must tell you its `sessionId`, which its own `fleet_whoami` gives it. Coordination only, **never** a task |
|
||||
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` reports a `collaborators` array, and each row carries that peer's `name` and the `sessionId` you send to. It is visible to you, to an architect and to another collaborator, never to a worker. Coordination only, **never** a task |
|
||||
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
|
||||
| Read your own held lead-to-lead mail (no ack) | `fleet_poll{coordId: <your own coord-id, from fleet_list's coordinator.selfId>}` — primary-only; never acks, so `fleet_list`'s `held[]` still shows it after. `fleet_list`'s `held[]` gives only a truncated preview — this is the only way to read the full body |
|
||||
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
|
||||
@@ -257,8 +257,8 @@ of those is refused at the gate, not queued.
|
||||
Two limits worth knowing before you hit them. **You cannot reach a worker** — not even to help one —
|
||||
because a worker belongs to the lead that spawned it, and routing around that would make you a
|
||||
second orchestrator with no plan. Send to the lead instead. And **you cannot read a ticket**, so you
|
||||
cannot collect a delegation's reply; ticket ids are a plain counter with no owner check, so holding
|
||||
one would let you walk every other session's answers.
|
||||
cannot collect a delegation's reply: `fleet_poll` refuses you at the role gate, and a ticket also
|
||||
records the terminal that created it, so even a leaked id reads nothing.
|
||||
|
||||
Being named buys you a channel, not authority. Your `fleet_send` to a lead is coordination between
|
||||
peers: the lead owes you no obedience, and you owe it none.
|
||||
|
||||
@@ -75,6 +75,9 @@ public final class Authz {
|
||||
* {@code SEND} is refused, as if no terminal were a configured lead or collaborator — the
|
||||
* same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR} gives explicitly. Every other action's
|
||||
* result is identical to the four-argument form's, since none of them consult the classifier.
|
||||
*
|
||||
* <p>Its default classifier denies every collaborator, so a caller enforcing authorization
|
||||
* must use the four-argument form instead.
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession) {
|
||||
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR);
|
||||
|
||||
@@ -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<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> 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);
|
||||
}
|
||||
|
||||
@@ -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 <em>blocking</em> {@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
|
||||
* <em>blocking</em> {@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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -874,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());
|
||||
@@ -894,7 +897,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,60 @@ 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);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@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) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -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<MessageService.AskResult> 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<MessageService.Reply> 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
|
||||
|
||||
@@ -0,0 +1,280 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* Pins that no file under {@code src/main/java} calls the fail-open
|
||||
* {@link MessageService#poll(String)} overload. That overload skips the ownership check in
|
||||
* {@code MessageService}'s {@code ownsTicket} entirely, so a caller of it can read any session's
|
||||
* ticket. Every production caller must go through {@link MessageService#poll(String, String)}
|
||||
* and pass a {@code callerTerminal} explicitly, even when it is {@code null}.
|
||||
*
|
||||
* <p>This reads each file's own source text rather than reflecting on compiled bytecode, because
|
||||
* the risk is a future one-word edit at a call site, not a missing overload.
|
||||
*
|
||||
* <p>The scan below finds a violation by its receiver, {@code messages.poll(}, rather than the
|
||||
* bare method name, so it does not mistake {@link java.util.Queue#poll()} for a violation. That
|
||||
* anchor only covers a {@code MessageService} reached through a variable or field named
|
||||
* {@code messages}, so {@link #everyMessageServiceDeclarationIsNamedMessages} pins the naming
|
||||
* convention the anchor depends on: a declaration under any other name would be invisible to the
|
||||
* scan above, and must turn this second check red instead of passing silently.
|
||||
*/
|
||||
class MessageServicePollUsageTest {
|
||||
|
||||
private static final Path PRODUCTION_SOURCE = Path.of("src/main/java");
|
||||
|
||||
@Test
|
||||
void noProductionFileCallsTheSingleArgumentPollOverload() throws IOException {
|
||||
List<String> violations = new ArrayList<>();
|
||||
List<String> twoArgSites = new ArrayList<>();
|
||||
int filesScanned = scanForPollCalls(PRODUCTION_SOURCE, violations, twoArgSites);
|
||||
|
||||
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "no
|
||||
// violations found" having looked at nothing.
|
||||
assertTrue(filesScanned > 0, "control failed: the scan under " + PRODUCTION_SOURCE
|
||||
+ " visited zero .java files -- the path is wrong, so the absence of violations "
|
||||
+ "below proves nothing");
|
||||
|
||||
assertTrue(violations.isEmpty(), "found a call to the fail-open MessageService.poll(String) "
|
||||
+ "overload, which skips the ownership check entirely -- pass a callerTerminal "
|
||||
+ "explicitly (even if null) through poll(String, String) instead: " + violations);
|
||||
|
||||
// CONTROL: the arity parser actually finds the two genuine two-argument call sites (the
|
||||
// MCP handler in FleetMcp and the REST handler in FleetApp). If this drops, the parser
|
||||
// itself is broken, not the production code -- a broken parser (or a scan root that
|
||||
// reaches no real source) must fail loudly here rather than pass vacuously above.
|
||||
assertEquals(2, twoArgSites.size(), "control failed: expected exactly the two known "
|
||||
+ "two-argument messages.poll(...) call sites, found: " + twoArgSites);
|
||||
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetMcp.java")),
|
||||
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerTerminal) "
|
||||
+ "site among: " + twoArgSites);
|
||||
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetApp.java")),
|
||||
"control failed: did not find the FleetApp.java messages.poll(...) site among: "
|
||||
+ twoArgSites);
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@code messages.poll(} anchor above only sees a {@code MessageService} reached through
|
||||
* a variable, field, or parameter named {@code messages}. This asserts that every such
|
||||
* declaration under {@code src/main/java} uses that name, so a differently named declaration
|
||||
* -- invisible to the scan above -- fails loudly here instead of letting that scan pass on a
|
||||
* call site it never looked at.
|
||||
*/
|
||||
@Test
|
||||
void everyMessageServiceDeclarationIsNamedMessages() throws IOException {
|
||||
List<String> names = new ArrayList<>();
|
||||
int filesScanned = scanForDeclarationNames(PRODUCTION_SOURCE, names);
|
||||
|
||||
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "every
|
||||
// declaration is named messages" having looked at nothing.
|
||||
assertTrue(filesScanned > 0, "control failed: the scan under " + PRODUCTION_SOURCE
|
||||
+ " visited zero .java files -- the path is wrong, so the result below proves nothing");
|
||||
|
||||
// CONTROL: the declaration pattern actually finds real declarations. Zero means the
|
||||
// pattern is broken, not that every MessageService variable, field, or parameter vanished.
|
||||
assertTrue(names.size() > 0, "control failed: found zero MessageService declarations under "
|
||||
+ PRODUCTION_SOURCE + " -- the declaration pattern is broken, update it before "
|
||||
+ "trusting the naming check below");
|
||||
|
||||
List<String> other = names.stream().filter(n -> !n.equals("messages")).distinct().toList();
|
||||
assertTrue(other.isEmpty(), "found a MessageService declaration not named \"messages\": "
|
||||
+ other + " -- the messages.poll( scan above only looks for that name, so a call "
|
||||
+ "through a differently named variable or field is invisible to it; either rename "
|
||||
+ "the declaration or widen that scan's anchor to cover it");
|
||||
}
|
||||
|
||||
private static int scanForPollCalls(Path root, List<String> violations, List<String> twoArgSites)
|
||||
throws IOException {
|
||||
List<Path> files = javaFiles(root);
|
||||
for (Path file : files) {
|
||||
scanFileForPollCalls(file, violations, twoArgSites);
|
||||
}
|
||||
return files.size();
|
||||
}
|
||||
|
||||
private static int scanForDeclarationNames(Path root, List<String> names) throws IOException {
|
||||
List<Path> files = javaFiles(root);
|
||||
Pattern declaration = Pattern.compile("MessageService\\s+([A-Za-z_][A-Za-z0-9_]*)");
|
||||
for (Path file : files) {
|
||||
if (file.getFileName().toString().equals("MessageService.java")) {
|
||||
continue; // the type's own declaration, not a caller holding a reference to it
|
||||
}
|
||||
String stripped = stripComments(Files.readString(file));
|
||||
Matcher m = declaration.matcher(stripped);
|
||||
while (m.find()) {
|
||||
int j = m.end();
|
||||
while (j < stripped.length() && Character.isWhitespace(stripped.charAt(j))) j++;
|
||||
if (j < stripped.length() && stripped.charAt(j) == '(') {
|
||||
continue; // a method named like the convention, e.g. "MessageService messages()"
|
||||
}
|
||||
names.add(m.group(1));
|
||||
}
|
||||
}
|
||||
return files.size();
|
||||
}
|
||||
|
||||
private static List<Path> javaFiles(Path root) throws IOException {
|
||||
try (Stream<Path> paths = Files.walk(root)) {
|
||||
return paths.filter(p -> p.toString().endsWith(".java")).toList();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Replaces {@code //} and {@code /* *}{@code /} comment text with nothing, leaving code,
|
||||
* string/char literals and line breaks untouched -- so a comment that merely mentions
|
||||
* {@code MessageService} in prose can never be read as a declaration.
|
||||
*/
|
||||
private static String stripComments(String source) {
|
||||
StringBuilder out = new StringBuilder(source.length());
|
||||
boolean inString = false;
|
||||
boolean inChar = false;
|
||||
int i = 0;
|
||||
while (i < source.length()) {
|
||||
char c = source.charAt(i);
|
||||
if (inString) {
|
||||
out.append(c);
|
||||
if (c == '\\' && i + 1 < source.length()) { out.append(source.charAt(i + 1)); i += 2; continue; }
|
||||
if (c == '"') inString = false;
|
||||
i++;
|
||||
continue;
|
||||
}
|
||||
if (inChar) {
|
||||
out.append(c);
|
||||
if (c == '\\' && i + 1 < source.length()) { out.append(source.charAt(i + 1)); i += 2; continue; }
|
||||
if (c == '\'') inChar = false;
|
||||
i++;
|
||||
continue;
|
||||
}
|
||||
if (c == '"') { inString = true; out.append(c); i++; continue; }
|
||||
if (c == '\'') { inChar = true; out.append(c); i++; continue; }
|
||||
if (c == '/' && i + 1 < source.length() && source.charAt(i + 1) == '/') {
|
||||
while (i < source.length() && source.charAt(i) != '\n') i++;
|
||||
continue; // leaves the newline itself for the next iteration to append
|
||||
}
|
||||
if (c == '/' && i + 1 < source.length() && source.charAt(i + 1) == '*') {
|
||||
i += 2;
|
||||
while (i < source.length() && !(source.charAt(i) == '*' && i + 1 < source.length()
|
||||
&& source.charAt(i + 1) == '/')) {
|
||||
if (source.charAt(i) == '\n') out.append('\n');
|
||||
i++;
|
||||
}
|
||||
i += 2;
|
||||
continue;
|
||||
}
|
||||
out.append(c);
|
||||
i++;
|
||||
}
|
||||
return out.toString();
|
||||
}
|
||||
|
||||
private static void scanFileForPollCalls(Path file, List<String> violations, List<String> twoArgSites)
|
||||
throws IOException {
|
||||
String source = Files.readString(file);
|
||||
String needle = "messages.poll(";
|
||||
int from = 0;
|
||||
int idx;
|
||||
while ((idx = source.indexOf(needle, from)) >= 0) {
|
||||
int argsStart = idx + needle.length();
|
||||
String args = extractBalancedArgs(source, argsStart, file, idx);
|
||||
int closeParenIndex = argsStart + args.length();
|
||||
from = closeParenIndex + 1;
|
||||
|
||||
if (args.isBlank()) {
|
||||
continue; // MessageService has no zero-argument poll() -- nothing to classify
|
||||
}
|
||||
String site = file + ":" + lineOf(source, idx);
|
||||
if (topLevelCommaCount(args) == 0) {
|
||||
violations.add(site + " -- messages.poll(" + args.trim() + ")");
|
||||
} else {
|
||||
twoArgSites.add(site);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The text between {@code messages.poll(} and its matching close paren: balanced over nested
|
||||
* calls, and never split by a paren or comma sitting inside a string or char literal.
|
||||
*/
|
||||
private static String extractBalancedArgs(String source, int start, Path file, int callIndex) {
|
||||
int depth = 1;
|
||||
boolean inString = false;
|
||||
boolean inChar = false;
|
||||
int i = start;
|
||||
while (i < source.length()) {
|
||||
char c = source.charAt(i);
|
||||
if (inString) {
|
||||
if (c == '\\') { i += 2; continue; }
|
||||
if (c == '"') inString = false;
|
||||
} else if (inChar) {
|
||||
if (c == '\\') { i += 2; continue; }
|
||||
if (c == '\'') inChar = false;
|
||||
} else if (c == '"') {
|
||||
inString = true;
|
||||
} else if (c == '\'') {
|
||||
inChar = true;
|
||||
} else if (c == '(') {
|
||||
depth++;
|
||||
} else if (c == ')') {
|
||||
depth--;
|
||||
if (depth == 0) return source.substring(start, i);
|
||||
}
|
||||
i++;
|
||||
}
|
||||
throw new IllegalStateException(
|
||||
"unbalanced parentheses scanning " + file + ":" + lineOf(source, callIndex));
|
||||
}
|
||||
|
||||
/**
|
||||
* Commas at paren/bracket/brace depth zero, skipping string and char literals -- the argument
|
||||
* separators a human reader would see, not every comma character in the text.
|
||||
*/
|
||||
private static int topLevelCommaCount(String args) {
|
||||
int depth = 0;
|
||||
int commas = 0;
|
||||
boolean inString = false;
|
||||
boolean inChar = false;
|
||||
int i = 0;
|
||||
while (i < args.length()) {
|
||||
char c = args.charAt(i);
|
||||
if (inString) {
|
||||
if (c == '\\') { i += 2; continue; }
|
||||
if (c == '"') inString = false;
|
||||
} else if (inChar) {
|
||||
if (c == '\\') { i += 2; continue; }
|
||||
if (c == '\'') inChar = false;
|
||||
} else if (c == '"') {
|
||||
inString = true;
|
||||
} else if (c == '\'') {
|
||||
inChar = true;
|
||||
} else if (c == '(' || c == '[' || c == '{') {
|
||||
depth++;
|
||||
} else if (c == ')' || c == ']' || c == '}') {
|
||||
depth--;
|
||||
} else if (c == ',' && depth == 0) {
|
||||
commas++;
|
||||
}
|
||||
i++;
|
||||
}
|
||||
return commas;
|
||||
}
|
||||
|
||||
private static int lineOf(String source, int index) {
|
||||
int line = 1;
|
||||
for (int i = 0; i < index; i++) {
|
||||
if (source.charAt(i) == '\n') line++;
|
||||
}
|
||||
return line;
|
||||
}
|
||||
}
|
||||
@@ -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<MessageService.Reply> 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<MessageService.AskResult> 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<MessageService.Reply> 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<MessageService.AskResult> 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<MessageService.Reply> 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());
|
||||
|
||||
@@ -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;
|
||||
@@ -139,6 +142,257 @@ 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 /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
|
||||
* 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();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@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<MessageService.AskResult> 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<MessageService.Reply> 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
|
||||
* 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