Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2dff4b84a4 |
@@ -0,0 +1,96 @@
|
||||
# Audit: async ticket / rendezvous lifecycle (`fleetd/src/main/java/dev/ltms/fleet/msg/`)
|
||||
|
||||
Scope: `Rendezvous.java` and `MessageService.java` — the lifecycle of an async ticket
|
||||
(`Task`) and a rendezvous waiter: create, send, ask, answer, resolve, timeout, abandon, prune.
|
||||
|
||||
## Main finding
|
||||
|
||||
```
|
||||
1. fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java:938
|
||||
2. issue: answer() completes an async ticket's future with a QUESTION outcome when the worker
|
||||
asks a second fleet_ask in the same resumed turn, permanently mislabeling a live delegation
|
||||
as failed and losing its real reply to fleet_poll.
|
||||
3. fix: guard the finishAsyncTask(turnId, result) call at line 938 the same way sendAsync's
|
||||
lambda already guards its own call (lines 999-1005): skip it when result.outcome() ==
|
||||
Outcome.QUESTION, and instead re-associate the task with the new turnId (as markAsyncQuestion
|
||||
does on the first ask).
|
||||
4. severity: high
|
||||
```
|
||||
|
||||
### Call sequence that reaches it
|
||||
|
||||
1. Lead: `fleet_send{sessionId: W, content: "task", wait:false}` → `sendAsync` creates `task1`
|
||||
/ `ticket1`. Its worker thread calls `send(W, content, ASYNC_TIMEOUT_MS, onAccepted, task1)`,
|
||||
which does `asyncTasksByWaiter.put(reply, task1)` (line 802) before blocking on
|
||||
`reply.get()`.
|
||||
2. Worker `W` calls `fleet_ask{"Q1"}` → `ask(W, "Q1", t)`. `markAsyncQuestion` finds `task1`
|
||||
via `asyncTasksByWaiter`, stamps `task1.turnId = turnId1`,
|
||||
`asyncTasksByTurn[turnId1] = task1`. `resolveQuestion` wakes step 1's `send()`, which returns
|
||||
`Outcome.QUESTION`; `sendAsync`'s lambda sees `QUESTION` and deliberately does **not** call
|
||||
`finishAsyncTask` (lines 1000-1005) — `ticket1` correctly polls `Phase.ASKING`.
|
||||
3. Lead polls, sees `ASKING`, answers: `fleet_send{turnId: turnId1, content: "A1"}` →
|
||||
`answer(turnId1, "A1", t)`. This opens a **new** waiter via `rendezvous.open(workerSession)`
|
||||
(line 929) but — unlike `send()` — never puts it into `asyncTasksByWaiter`.
|
||||
`answerAsk(turnId1, "A1")` unblocks the worker's `ask()` call.
|
||||
`clearAsyncQuestion(turnId1, false)` clears `task1.question` but keeps
|
||||
`asyncTasksByTurn[turnId1] = task1` (deliberate, per its own javadoc). `answer()` then blocks
|
||||
on its own `reply.get()` (line 936).
|
||||
4. Worker `W`, still in the same resumed turn, calls `fleet_ask{"Q2"}` again before replying →
|
||||
a second `ask(W, "Q2", t)`. `openAsk` mints `turnId2`. `markAsyncQuestion` looks up
|
||||
`asyncTasksByWaiter.get(waiter)` for the waiter `answer()` opened in step 3 — **not found**
|
||||
(never registered), so `task == null`; `task1.turnId` stays `turnId1`, no
|
||||
`asyncTasksByTurn[turnId2]` entry is ever created. `resolveQuestion` still succeeds (it only
|
||||
needs a live waiter, not a `Task`) and wakes `answer(turnId1,...)`'s blocked `reply.get()`
|
||||
with `Resolution(QUESTION, "Q2", turnId2)`.
|
||||
5. `answer(turnId1,...)` (line 936-939): `result = Reply(Outcome.QUESTION, "Q2", turnId2)`;
|
||||
`finishAsyncTask(turnId1, result)` looks up `task1` by the **original** `turnId1` (still
|
||||
stamped from step 3) and unconditionally does `task1.future.complete(result)` — completing
|
||||
`ticket1`'s future with a **QUESTION** outcome, then removes `asyncTasksByTurn[turnId1]`.
|
||||
`answer()` returns `Outcome.QUESTION` to the lead's own `fleet_send{turnId1,...}` call
|
||||
(correct, and separately answerable via `turnId2`), but `ticket1` is now terminally done.
|
||||
|
||||
### What goes wrong
|
||||
|
||||
- `fleet_poll{ticket1}` now hits the `f.isDone()` branch in `poll()` permanently. `Outcome.QUESTION`
|
||||
is not `REPLIED`/`COMPLETED_UNREPLIED` (`r.completed()` is false) and carries no
|
||||
`WORKER_FAILED`/`BACKEND_EXHAUSTED` reason, so it falls through to
|
||||
`Phase.FAILED`, `detail = "no reply — question"` — even though the worker is alive and only
|
||||
waiting on `turnId2`.
|
||||
- `asyncTasksByTurn` no longer has any entry for `task1`/`target`, so
|
||||
`hasAsyncQuestion(target)` goes back to `false` immediately, and `hasOrphanedDelegation`
|
||||
no longer excludes this target's real state correctly either.
|
||||
- If the worker's eventual real `fleet_reply` (after `turnId2` is answered, or times out and it
|
||||
finishes on its own) is not captured by a chained direct `answer(turnId2,...)` call,
|
||||
`reply()`'s fast path (`rendezvous.resolve`) finds no live waiter, `askAnsweredAsyncTasks`
|
||||
finds no candidate (`task1.future.isDone()` is already true, so it is excluded), and the reply
|
||||
is silently dropped into the inbox as a **stranded reply** — unreachable from `ticket1` and
|
||||
from `abandon()`'s stranded-reply recovery (no open `matching` task exists any more).
|
||||
|
||||
### Confidence
|
||||
|
||||
High. I traced this with no races or interleavings assumed beyond the documented, deterministic
|
||||
CB-205 chained-ask protocol that `outcomeOf`/`answer()` already generically support (mapping
|
||||
`Kind.QUESTION` through `answer()`'s own return value is clearly intentional — see the
|
||||
`Reply.turnId()` javadoc). I did not run the suite, but grepped
|
||||
`src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` for a test exercising a *second*
|
||||
`fleet_ask` inside one resumed (answered) turn and found none — `asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter`
|
||||
and `askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn` both cover only a single ask per
|
||||
turn. A `git log -p` on this file also turned up the CB-588 comment (now at
|
||||
`sendAsync`, describing the `whenComplete` hook) which explicitly frames
|
||||
`answer()`'s `finishAsyncTask(turnId, result)` as firing "once a QUESTION is resolved" —
|
||||
i.e. the author modeled that call as inherently terminal, which is exactly the assumption this
|
||||
bug violates when the resumed turn asks again.
|
||||
|
||||
## Secondary (much shorter)
|
||||
|
||||
1. **Root-cause detail, same defect as above** — `answer()` (line 929) never puts its freshly
|
||||
opened waiter into `asyncTasksByWaiter`, unlike `send()` (line 802). Even if the outcome
|
||||
guard above is added, a chained second ask still can't be re-attached to `task1` via the
|
||||
normal `markAsyncQuestion` path without also fixing this registration gap.
|
||||
2. **Low / shape only** — `pruneTerminalTickets()` (line ~1092) is only invoked from inside
|
||||
`sendAsync()`. A fleet whose sessions stop receiving new async sends (e.g. everything now
|
||||
goes through blocking `send()`, or the target churns and workers are torn down) never prunes
|
||||
its already-terminal `tasks` entries past `TICKET_TTL_NANOS`. Not reachable as a "stuck"
|
||||
ticket (tickets still resolve correctly), only as unbounded `tasks`/`asyncTasksByWaiter`-adjacent
|
||||
memory growth over a long-lived daemon with no further `sendAsync` traffic; did not verify
|
||||
this is realistic in production traffic patterns, flagging as a shape only.
|
||||
@@ -257,7 +257,7 @@ public final class FleetMcp {
|
||||
// Each handler is built once and wired to its fleet_* tool below.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> sendHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_send", req.arguments()),
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SEND,
|
||||
str(req.arguments(), "sessionId"));
|
||||
if (denied != null) return denied;
|
||||
String caller = callerTerminal(exchange);
|
||||
@@ -299,7 +299,7 @@ public final class FleetMcp {
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> replyHandler =
|
||||
(exchange, req) -> {
|
||||
String self = callerTerminal(exchange);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_reply", req.arguments()), self);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.REPLY, self);
|
||||
if (denied != null) return denied;
|
||||
return reply(messages, self, str(req.arguments(), "content"));
|
||||
};
|
||||
@@ -307,13 +307,13 @@ public final class FleetMcp {
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> askHandler =
|
||||
(exchange, req) -> {
|
||||
String self = callerTerminal(exchange);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_ask", req.arguments()), self);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.ASK, self);
|
||||
if (denied != null) return denied;
|
||||
return ask(messages, self, str(req.arguments(), "question"), timeoutMs(req.arguments()));
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> statusHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return status(messages, str(req.arguments(), "sessionId"));
|
||||
};
|
||||
@@ -322,7 +322,7 @@ public final class FleetMcp {
|
||||
Map<String, Object> a = req.arguments();
|
||||
String target = str(a, "target");
|
||||
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
|
||||
McpSchema.CallToolResult denied = deny(exchange, pollAction(target), target);
|
||||
if (denied != null) return denied;
|
||||
return poll(messages, str(a, "ticket"), target);
|
||||
};
|
||||
@@ -331,14 +331,14 @@ public final class FleetMcp {
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> ackHandler =
|
||||
(exchange, req) -> {
|
||||
Map<String, Object> a = req.arguments();
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_ack", a), str(a, "target"));
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.DRAIN, str(a, "target"));
|
||||
if (denied != null) return denied;
|
||||
return ack(messages, str(a, "target"), str(a, "msgId"));
|
||||
};
|
||||
// Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> spawnHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_spawn", req.arguments()), null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null);
|
||||
if (denied != null) return denied;
|
||||
String caller = callerTerminal(exchange);
|
||||
// SPAWN is already auth-gated to PRIMARY (architects can never call it), but
|
||||
@@ -355,7 +355,7 @@ public final class FleetMcp {
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> listHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, callers == null ? Map.of() : callers.leads(),
|
||||
@@ -365,19 +365,19 @@ public final class FleetMcp {
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
String paneId = str(req.arguments(), "paneId");
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_stop", req.arguments()), paneId);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.STOP, paneId);
|
||||
if (denied != null) return denied;
|
||||
return stop(sessions, paneId);
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> profilesHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_profiles", Map.of()), null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return profiles(workers, quarantine, outage);
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> whoamiHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_whoami", Map.of()), null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return whoami(principal(exchange), sessions);
|
||||
};
|
||||
@@ -754,25 +754,6 @@ public final class FleetMcp {
|
||||
return isBlank(target) ? Authz.Action.READ : Authz.Action.DRAIN;
|
||||
}
|
||||
|
||||
/**
|
||||
* The action a registered tool handler actually hands to the authorization gate.
|
||||
* Keeping this choice beside the registered-tool inventory makes a new tool fail the coverage
|
||||
* test until its action is pinned.
|
||||
*/
|
||||
static Authz.Action toolAction(String toolName, Map<String, Object> arguments) {
|
||||
return switch (toolName) {
|
||||
case "fleet_send" -> Authz.Action.SEND;
|
||||
case "fleet_reply" -> Authz.Action.REPLY;
|
||||
case "fleet_ask" -> Authz.Action.ASK;
|
||||
case "fleet_status", "fleet_list", "fleet_profiles", "fleet_whoami" -> Authz.Action.READ;
|
||||
case "fleet_poll" -> pollAction(str(arguments, "target"));
|
||||
case "fleet_ack" -> Authz.Action.DRAIN;
|
||||
case "fleet_spawn" -> Authz.Action.SPAWN;
|
||||
case "fleet_stop" -> Authz.Action.STOP;
|
||||
default -> throw new IllegalArgumentException("unregistered tool: " + toolName);
|
||||
};
|
||||
}
|
||||
|
||||
/** {@code fleet_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
|
||||
static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
|
||||
if (!isBlank(target)) {
|
||||
|
||||
@@ -47,22 +47,6 @@ import java.util.stream.Collectors;
|
||||
*/
|
||||
public final class FleetApp {
|
||||
|
||||
/** The authorization action the matching route handler hands to {@link #allow}. */
|
||||
static Authz.Action routeAction(String route) {
|
||||
return switch (route) {
|
||||
case "GET /metrics" -> Authz.Action.METRICS;
|
||||
case "POST /members" -> Authz.Action.SPAWN;
|
||||
case "DELETE /members/{paneId}" -> Authz.Action.STOP;
|
||||
case "POST /sessions/{id}/message" -> Authz.Action.SEND;
|
||||
case "POST /sessions/{id}/reply" -> Authz.Action.REPLY;
|
||||
case "GET /sessions/{id}/replies" -> Authz.Action.DRAIN;
|
||||
case "POST /sessions/{id}/ask" -> Authz.Action.ASK;
|
||||
case "GET /sessions", "GET /agents", "GET /members", "GET /profiles",
|
||||
"GET /member-credentials", "GET /sessions/{id}/status", "GET /tasks/{ticket}" -> Authz.Action.READ;
|
||||
default -> throw new IllegalArgumentException("route has no authorization gate: " + route);
|
||||
};
|
||||
}
|
||||
|
||||
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
|
||||
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
|
||||
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_000;
|
||||
@@ -238,7 +222,7 @@ public final class FleetApp {
|
||||
|
||||
/** Prometheus scrape endpoint (CB-502). */
|
||||
private void metrics(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /metrics"), null)) {
|
||||
if (!allow(ctx, Authz.Action.METRICS, null)) {
|
||||
return;
|
||||
}
|
||||
ctx.status(200).contentType("text/plain; version=0.0.4; charset=utf-8").result(metrics.render());
|
||||
@@ -310,7 +294,7 @@ public final class FleetApp {
|
||||
* member workspace (they live on the member daemon only).
|
||||
*/
|
||||
private void sessions(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /sessions"), null)) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
return;
|
||||
}
|
||||
List<Map<String, Object>> out = new ArrayList<>();
|
||||
@@ -335,7 +319,7 @@ public final class FleetApp {
|
||||
|
||||
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
|
||||
private void agents(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /agents"), null)) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
return;
|
||||
}
|
||||
ctx.status(200).json(Map.of("agents",
|
||||
@@ -344,7 +328,7 @@ public final class FleetApp {
|
||||
|
||||
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
|
||||
private void listMembers(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /members"), null)) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
return;
|
||||
}
|
||||
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
|
||||
@@ -375,7 +359,7 @@ public final class FleetApp {
|
||||
|
||||
/** The configured worker profiles and which one a no-argument spawn uses. */
|
||||
private void profiles(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /profiles"), null)) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
return;
|
||||
}
|
||||
ctx.status(200).json(Map.of(
|
||||
@@ -392,7 +376,7 @@ public final class FleetApp {
|
||||
* name list, which is exactly what let the list drift silently behind the real policy.
|
||||
*/
|
||||
private void memberCredentials(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /member-credentials"), null)) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
return;
|
||||
}
|
||||
MemberCredentialPolicyView view = memberCredentials.get();
|
||||
@@ -412,7 +396,7 @@ public final class FleetApp {
|
||||
* the subscription boundary, 400 for an unknown profile.
|
||||
*/
|
||||
private void spawnMember(Context ctx) {
|
||||
if (!allow(ctx, routeAction("POST /members"), null)) {
|
||||
if (!allow(ctx, Authz.Action.SPAWN, null)) {
|
||||
return;
|
||||
}
|
||||
String role = ctx.queryParam("role");
|
||||
@@ -483,7 +467,7 @@ public final class FleetApp {
|
||||
/** Tear a worker down by pane id. */
|
||||
private void stopMember(Context ctx) {
|
||||
String paneId = ctx.pathParam("paneId");
|
||||
if (!allow(ctx, routeAction("DELETE /members/{paneId}"), paneId)) {
|
||||
if (!allow(ctx, Authz.Action.STOP, paneId)) {
|
||||
return;
|
||||
}
|
||||
sessions.release(paneId);
|
||||
@@ -498,7 +482,7 @@ public final class FleetApp {
|
||||
*/
|
||||
private void sendMessage(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, routeAction("POST /sessions/{id}/message"), id)) {
|
||||
if (!allow(ctx, Authz.Action.SEND, id)) {
|
||||
return;
|
||||
}
|
||||
String content;
|
||||
@@ -585,7 +569,7 @@ public final class FleetApp {
|
||||
*/
|
||||
private void askMessage(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, routeAction("POST /sessions/{id}/ask"), id)) {
|
||||
if (!allow(ctx, Authz.Action.ASK, id)) {
|
||||
return;
|
||||
}
|
||||
String question;
|
||||
@@ -624,7 +608,7 @@ public final class FleetApp {
|
||||
// The rule that matters: a worker may reply only as itself. Over MCP this was already true
|
||||
// structurally (identity comes from the connection, never an argument); over REST the path
|
||||
// id was simply trusted, so this is where the invariant actually gets enforced.
|
||||
if (!allow(ctx, routeAction("POST /sessions/{id}/reply"), id)) {
|
||||
if (!allow(ctx, Authz.Action.REPLY, id)) {
|
||||
return;
|
||||
}
|
||||
String content;
|
||||
@@ -645,7 +629,7 @@ public final class FleetApp {
|
||||
*/
|
||||
private void drainReplies(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, routeAction("GET /sessions/{id}/replies"), id)) {
|
||||
if (!allow(ctx, Authz.Action.DRAIN, id)) {
|
||||
return;
|
||||
}
|
||||
var replies = messages.drainReplies(id);
|
||||
@@ -663,7 +647,7 @@ public final class FleetApp {
|
||||
*/
|
||||
private void sessionStatus(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, routeAction("GET /sessions/{id}/status"), id)) {
|
||||
if (!allow(ctx, Authz.Action.READ, id)) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
@@ -688,7 +672,7 @@ public final class FleetApp {
|
||||
|
||||
/** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */
|
||||
private void taskStatus(Context ctx) {
|
||||
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
return;
|
||||
}
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"));
|
||||
|
||||
@@ -24,13 +24,8 @@ import io.modelcontextprotocol.spec.McpSchema;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -49,9 +44,6 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
*/
|
||||
class FleetMcpAuthzTest {
|
||||
|
||||
private static final Path MCP_SOURCE = Path.of("src/main/java/dev/ltms/fleet/mcp/FleetMcp.java");
|
||||
private static final Pattern TOOL_REGISTRATION = Pattern.compile("tool\\(\\\"(fleet_[a-z_]+)\\\"");
|
||||
|
||||
private final FakeHerdr herdr = new FakeHerdr();
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private Metrics metrics;
|
||||
@@ -209,42 +201,6 @@ class FleetMcpAuthzTest {
|
||||
"a blank target is an absent target");
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyRegisteredToolHasItsHandlerActionPinned() {
|
||||
Set<String> registered = toolsTheServerRegisters();
|
||||
assertTrue(registered.size() >= 10,
|
||||
"scraped only " + registered.size() + " tool registrations from FleetMcp (" + registered
|
||||
+ "); the server registers eleven, so the tool(\"…\") scrape has stopped matching");
|
||||
registered.forEach(tool -> assertDoesNotThrow(() -> FleetMcp.toolAction(tool, Map.of()),
|
||||
() -> tool + " is registered but has no pinned authorization action"));
|
||||
|
||||
assertEquals(Authz.Action.SEND, FleetMcp.toolAction("fleet_send", Map.of()));
|
||||
assertEquals(Authz.Action.REPLY, FleetMcp.toolAction("fleet_reply", Map.of()));
|
||||
assertEquals(Authz.Action.ASK, FleetMcp.toolAction("fleet_ask", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_status", Map.of()));
|
||||
assertEquals(Authz.Action.DRAIN, FleetMcp.toolAction("fleet_ack", Map.of()));
|
||||
assertEquals(Authz.Action.SPAWN, FleetMcp.toolAction("fleet_spawn", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_list", Map.of()));
|
||||
assertEquals(Authz.Action.STOP, FleetMcp.toolAction("fleet_stop", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_profiles", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_whoami", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_poll", Map.of("ticket", "task")));
|
||||
assertEquals(Authz.Action.DRAIN, FleetMcp.toolAction("fleet_poll", Map.of("target", "term_b")));
|
||||
}
|
||||
|
||||
private static Set<String> toolsTheServerRegisters() {
|
||||
try {
|
||||
Matcher matcher = TOOL_REGISTRATION.matcher(Files.readString(MCP_SOURCE));
|
||||
Set<String> tools = new LinkedHashSet<>();
|
||||
while (matcher.find()) {
|
||||
tools.add(matcher.group(1));
|
||||
}
|
||||
return tools;
|
||||
} catch (Exception e) {
|
||||
throw new AssertionError("could not scrape FleetMcp tool registrations", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aWorkerMayNotDrainAnotherSessionsInboxByPolling() {
|
||||
FleetMcp m = mcp(true);
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
package dev.ltms.fleet.rest;
|
||||
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
import dev.ltms.fleet.auth.Authz;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
@@ -26,15 +25,8 @@ import java.net.URI;
|
||||
import java.net.http.HttpClient;
|
||||
import java.net.http.HttpRequest;
|
||||
import java.net.http.HttpResponse;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -44,10 +36,6 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
*/
|
||||
class FleetAppAuthTest {
|
||||
|
||||
private static final Path REST_SOURCE = Path.of("src/main/java/dev/ltms/fleet/rest/FleetApp.java");
|
||||
private static final Pattern ROUTE_REGISTRATION =
|
||||
Pattern.compile("app\\.(get|post|delete|put|patch)\\(\\s*\"([^\"]+)\"");
|
||||
|
||||
private final HttpClient http = HttpClient.newHttpClient();
|
||||
private Javalin app;
|
||||
private Metrics metrics;
|
||||
@@ -104,49 +92,6 @@ class FleetAppAuthTest {
|
||||
return http.send(b.build(), HttpResponse.BodyHandlers.ofString());
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyRegisteredRouteHasItsHandlerActionPinned() {
|
||||
Set<String> registered = routesTheServerRegisters();
|
||||
assertTrue(registered.size() >= 15,
|
||||
"scraped only " + registered.size() + " route registrations from FleetApp (" + registered
|
||||
+ "); the app.<verb>(\"…\") scrape has stopped matching");
|
||||
// Liveness must work before credentials can be checked, so this route is deliberately open.
|
||||
assertTrue(registered.remove("GET /healthz"), "GET /healthz must stay an explicit ungated exception");
|
||||
registered.forEach(route -> assertDoesNotThrow(() -> FleetApp.routeAction(route),
|
||||
() -> route + " is registered but has no pinned authorization action"));
|
||||
assertEquals(Authz.Action.METRICS, FleetApp.routeAction("GET /metrics"));
|
||||
assertEquals(Authz.Action.SPAWN, FleetApp.routeAction("POST /members"));
|
||||
assertEquals(Authz.Action.STOP, FleetApp.routeAction("DELETE /members/{paneId}"));
|
||||
assertEquals(Authz.Action.SEND, FleetApp.routeAction("POST /sessions/{id}/message"));
|
||||
assertEquals(Authz.Action.REPLY, FleetApp.routeAction("POST /sessions/{id}/reply"));
|
||||
assertEquals(Authz.Action.DRAIN, FleetApp.routeAction("GET /sessions/{id}/replies"));
|
||||
assertEquals(Authz.Action.ASK, FleetApp.routeAction("POST /sessions/{id}/ask"));
|
||||
for (String route : Set.of("GET /sessions", "GET /agents", "GET /members", "GET /profiles",
|
||||
"GET /member-credentials", "GET /sessions/{id}/status", "GET /tasks/{ticket}")) {
|
||||
assertEquals(Authz.Action.READ, FleetApp.routeAction(route), route);
|
||||
}
|
||||
assertThrows(IllegalArgumentException.class, () -> FleetApp.routeAction("GET /healthz"));
|
||||
}
|
||||
|
||||
private static Set<String> routesTheServerRegisters() {
|
||||
try {
|
||||
String source = Files.readString(REST_SOURCE).lines()
|
||||
.filter(line -> {
|
||||
String stripped = line.stripLeading();
|
||||
return !(stripped.startsWith("//") || stripped.startsWith("*") || stripped.startsWith("/*"));
|
||||
})
|
||||
.collect(Collectors.joining("\n"));
|
||||
Matcher matcher = ROUTE_REGISTRATION.matcher(source);
|
||||
Set<String> routes = new LinkedHashSet<>();
|
||||
while (matcher.find()) {
|
||||
routes.add(matcher.group(1).toUpperCase(Locale.ROOT) + " " + matcher.group(2));
|
||||
}
|
||||
return routes;
|
||||
} catch (Exception e) {
|
||||
throw new AssertionError("could not scrape FleetApp route registrations", e);
|
||||
}
|
||||
}
|
||||
|
||||
// --- loopback-trust: the caller is the primary -------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user