Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 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);
|
||||
|
||||
@@ -694,7 +694,8 @@ public final class FleetApp {
|
||||
|
||||
if (!wait) {
|
||||
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
|
||||
String ticket = messages.sendAsync(id, content);
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
String ticket = messages.sendAsync(id, content, null, caller == null ? null : caller.terminal());
|
||||
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
|
||||
return;
|
||||
}
|
||||
@@ -894,7 +895,8 @@ public final class FleetApp {
|
||||
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
|
||||
return;
|
||||
}
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"));
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.terminal());
|
||||
if (v == null) {
|
||||
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
|
||||
return;
|
||||
|
||||
@@ -503,6 +503,33 @@ class FleetMcpAuthzTest {
|
||||
+ "not pass a literal boolean -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code fleet_poll{ticket}} must thread the calling connection's own terminal into
|
||||
* {@link MessageService#poll(String, String)}, so a worker cannot read a ticket a different
|
||||
* session created.
|
||||
*/
|
||||
@Test
|
||||
void theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
int start = source.indexOf("pollHandler =");
|
||||
assertTrue(start >= 0, "could not find the fleet_poll handler (pollHandler) in " + MCP_SOURCE
|
||||
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
|
||||
int end = source.indexOf("ackHandler =", start);
|
||||
assertTrue(end > start, "could not find the handler declared after pollHandler to bound the scrape");
|
||||
String handlerBlock = source.substring(start, end);
|
||||
|
||||
// CONTROL: the block we scraped really does call poll(...) -- if this fails, the anchors
|
||||
// above moved and the assertions below would otherwise pass on nothing.
|
||||
assertTrue(handlerBlock.contains("poll(messages,"),
|
||||
"control failed: the scraped pollHandler block contains no poll(messages, ...) call "
|
||||
+ "at all -- the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
|
||||
"the fleet_poll handler must thread callerTerminal(exchange) into poll(...), not omit "
|
||||
+ "it or pass a literal null -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
// --- which action each tool hands the gate (fleetd #272) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -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