Compare commits
29 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7cf6075b79 | |||
| fbdcd709c9 | |||
| 8d3f10d291 | |||
| 38f4fd64ee | |||
| 3fc39b981d | |||
| 70328ca0f8 | |||
| efab9b8c49 | |||
| 2374de28e4 | |||
| dd18bd1f38 | |||
| efeffb4ab7 | |||
| ed4f4b08ad | |||
| 28a1f3d6f5 | |||
| b12d70716b | |||
| 0f5985b419 | |||
| 6e06058b07 | |||
| e5f4fb81ab | |||
| d28ab0968b | |||
| 9a64d42599 | |||
| f4e0ca41e6 | |||
| 29a2f97c25 | |||
| d2db8c7dc9 | |||
| 95311c6e8e | |||
| 10ab58e4fc | |||
| 0ba597e394 | |||
| c468953963 | |||
| 25d53e6ef7 | |||
| a9a37af957 | |||
| b205bcc2aa | |||
| 11eccc3a1b |
@@ -100,7 +100,11 @@ below are the procedure — run them in order, every task, not only the big ones
|
||||
5. **Collect** — `fleet_poll{ticket}` → `fleet_ack{target, msgId}`. Answer a worker's `fleet_ask`
|
||||
with `fleet_send{turnId, content}` — **not** `sessionId`. A worker gone quiet is diagnosed with
|
||||
`fleet_status`, never by reading its terminal; it also reports an open question and the `turnId`
|
||||
that answers it. **A worker's ask waits ~55 seconds, and no nudge makes that longer** — so never
|
||||
that answers it — but **only to the caller that created that delegation**, so a question raised
|
||||
under an architect's brief is invisible to you, and seeing none does not mean there is none.
|
||||
**Only that same creator can answer it.** A `turnId` you came by any other way is refused, so an
|
||||
architect's worker waits for that architect and not for you.
|
||||
**A worker's ask waits ~55 seconds, and no nudge makes that longer** — so never
|
||||
brief a worker to "ask me". Decide before you delegate, or give it an explicit default.
|
||||
**A correction cannot reach a busy member.** A `fleet_send` to a working member is *accepted* and
|
||||
returns a ticket, and is then never delivered — measured here three times in one session, and the
|
||||
@@ -165,10 +169,10 @@ you decide.
|
||||
| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) + `loopHealth` (`RUNNING`, `STALLED`, or `STOPPED` for `statusPoller` and `sessionReaper`) · one peer's state: `fleet_status{sessionId}` |
|
||||
| Delegate (blocking) | `fleet_send{sessionId, content}` |
|
||||
| Delegate (long task) | `fleet_send{sessionId, content, wait:false}` → ticket → `fleet_poll{ticket}` |
|
||||
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
|
||||
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId`, and only the caller that created that delegation |
|
||||
| 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 +261,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.
|
||||
@@ -330,10 +334,22 @@ must obey belongs in the charter, not here.
|
||||
(#362). **Read `plugin/` before designing anything about onboarding a project.** Two limits are
|
||||
structural, not bugs: a plugin cannot carry the role agent files, because
|
||||
`ClaudeCodeLauncher.java:371` requires `<cwd>/.claude/agents/<role>.md` in the member's own
|
||||
worktree; and a plugin cannot deliver anything to members at all, because
|
||||
`ClaudeCodeLauncher.java:285` exports `CLAUDE_CONFIG_DIR` and every Claude profile here sets it,
|
||||
so a member never reads the operator's plugin store. **The plugin is the lead-side surface;
|
||||
member-facing assets travel in the worktree.**
|
||||
worktree; and a plugin reaches a member only through `CLAUDE_CONFIG_DIR`, which
|
||||
`ClaudeCodeLauncher.java:286` exports with `putIfPresent` — so only for a profile that sets
|
||||
`configDir`. Every `claude-code` profile does set one (the four without are `opencode`, which
|
||||
never reads that variable). **But measured 2026-10-04: two of them point at
|
||||
`~/.ccs/instances/ltms`, which is the operator's own `CLAUDE_CONFIG_DIR` on this host.** So for an
|
||||
`opus` or `sonnet` member, "a member never reads the operator's plugin store" is false — it reads
|
||||
the same store, because that store is the one its `configDir` names. It stays true for `local` and
|
||||
`local-direct`, which point at `~/.ccs/instances/gx10`. `ClaudeCodeLauncher`'s own javadoc names
|
||||
the related hazard: that file is rewritten on every spawn, so for those two profiles fleetd and the
|
||||
operator's live session write the same `.claude.json`, and its compare-and-swap "narrows the
|
||||
lost-update window, it does not close it". Re-measure which profiles share the operator's dir with
|
||||
`awk '/^profiles:/{i=1;next} /^[a-z]/{i=0} i&&/^ [a-z-]+:$/{p=$1} i&&/configDir:/{print p,$2}'
|
||||
fleetd/fleetd.yaml` against `echo $CLAUDE_CONFIG_DIR`; delete this note once no profile names the
|
||||
operator's dir. **Treat the plugin as the lead-side surface and put member-facing assets in the
|
||||
worktree** — that conclusion holds either way, because a worktree asset does not depend on which
|
||||
config dir a member reads.
|
||||
- **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/`
|
||||
(a submodule with its own remote).
|
||||
- **A provisioned worktree neutralizes `.mcp.json`, `opencode.json` and `.autoenv`** — the repo's
|
||||
|
||||
@@ -51,7 +51,6 @@ import dev.ltms.fleet.power.CaffeinateSleepAssertionMechanism;
|
||||
import dev.ltms.fleet.power.IdleSleepGuard;
|
||||
import dev.ltms.fleet.rest.FleetApp;
|
||||
import dev.ltms.fleet.session.GitWorktrees;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.SessionReaper;
|
||||
import io.javalin.Javalin;
|
||||
@@ -480,13 +479,9 @@ final class FleetdAssembly {
|
||||
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||
|
||||
// fleetd #669 Unit D: a live spawned member resolves as its own role, whatever a tab map
|
||||
// says about the same terminal — read from the roster meant for a hot path (SessionManager
|
||||
// javadoc), never rosterResolved(), since resolve() runs on every request.
|
||||
Function<String, MemberRole> spawnedMemberRole = terminal -> sessions.roster().stream()
|
||||
.filter(s -> terminal.equals(s.terminalId()))
|
||||
.map(MemberSession::role)
|
||||
.findFirst()
|
||||
.orElse(null);
|
||||
// says about the same terminal. fleetd #702: SessionManager.spawnedMemberRole also answers
|
||||
// for a pane mid-teardown, not only one still in the registry — see its javadoc.
|
||||
Function<String, MemberRole> spawnedMemberRole = sessions::spawnedMemberRole;
|
||||
|
||||
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
|
||||
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -249,17 +249,6 @@ public final class CallerResolver {
|
||||
return leadTerminals.get();
|
||||
}
|
||||
|
||||
/**
|
||||
* The currently-recognised architect slots, {@code terminal_id → slot name} (CB-548).
|
||||
*
|
||||
* <p>Read from the same supplier {@link #resolve} consults, so a slot that is <em>listed</em>
|
||||
* here but would not <em>resolve</em> (or the reverse) cannot drift apart. Live for the same
|
||||
* reason as {@link #leads()}.
|
||||
*/
|
||||
public Map<String, String> members() {
|
||||
return architectTerminals.get();
|
||||
}
|
||||
|
||||
/**
|
||||
* The currently-recognised collaborator tabs, {@code terminal_id → name}.
|
||||
*
|
||||
|
||||
@@ -480,7 +480,7 @@ public final class FleetMcp {
|
||||
// Answering a worker's fleet_ask (CB-205): resolve its blocked question and
|
||||
// block for the worker's reply as it resumes the same turn. This is the same
|
||||
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
|
||||
return answer(messages, turnId, content, timeoutMs(a));
|
||||
return answer(messages, turnId, content, timeoutMs(a), caller);
|
||||
}
|
||||
// CB-548: delegator ownership (which lead's reply nudge this worker routes to,
|
||||
// CB-532) is recorded only once the send is ACCEPTED — MessageService has won the
|
||||
@@ -490,8 +490,8 @@ public final class FleetMcp {
|
||||
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
|
||||
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
|
||||
return Boolean.FALSE.equals(a.get("wait"))
|
||||
? sendAsync(messages, target, content, onAccepted, workers.profiles())
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
|
||||
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), caller);
|
||||
};
|
||||
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
|
||||
// is "is this caller a worker at all", and it can only ever reply as itself.
|
||||
@@ -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) -> {
|
||||
@@ -524,7 +524,7 @@ public final class FleetMcp {
|
||||
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
|
||||
if (denied != null) return denied;
|
||||
return poll(messages, leadChannel, str(a, "ticket"), target, coordId);
|
||||
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
|
||||
};
|
||||
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
|
||||
// Acking removes a reply from the inbox, so it is a drain, not a read.
|
||||
@@ -562,7 +562,8 @@ public final class FleetMcp {
|
||||
callerTerminal(exchange),
|
||||
callers.collaborators(), collaboratorsVisibleTo(principal(exchange)),
|
||||
new CoordinationSource(leadChannel, peers),
|
||||
coordinatorVisibleTo(principal(exchange)));
|
||||
coordinatorVisibleTo(principal(exchange)),
|
||||
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)));
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
@@ -759,7 +760,38 @@ public final class FleetMcp {
|
||||
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
|
||||
}
|
||||
|
||||
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
|
||||
/**
|
||||
* Who may see {@code fleet_list}'s {@code leads} array — exactly the roles that may
|
||||
* {@link Authz.Action#SEND} to a lead: the primary, an architect, and a collaborator. A
|
||||
* collaborator's own {@code fleet_whoami} carries no lead address, and {@code leads} is the
|
||||
* only place this tool gives one, so a collaborator needs this array to use the send it
|
||||
* already holds. A worker can never {@code SEND} at all, so it still sees neither this array
|
||||
* nor {@code members}; a worker's own facts come from {@code fleet_whoami} instead. Split
|
||||
* out for the same reason as {@link #coordinatorVisibleTo} and
|
||||
* {@link #collaboratorsVisibleTo}: the decision must be unit-testable without fabricating an
|
||||
* SDK {@code McpSyncServerExchange}, and the handler must call this named predicate rather
|
||||
* than inlining the check.
|
||||
*/
|
||||
static boolean leadsVisibleTo(Principal caller) {
|
||||
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
|
||||
}
|
||||
|
||||
/**
|
||||
* Who may see {@code fleet_list}'s {@code members} array — the primary and an architect,
|
||||
* which may {@link Authz.Action#SEND} to a member. A worker can never {@code SEND} at all,
|
||||
* and a collaborator may {@code SEND} only to a lead or another collaborator, never to a
|
||||
* spawned member, so neither sees this array even though {@link #leadsVisibleTo} grants a
|
||||
* collaborator the sibling one.
|
||||
*/
|
||||
static boolean membersVisibleTo(Principal caller) {
|
||||
return caller.isPrimary() || caller.isArchitect();
|
||||
}
|
||||
|
||||
/**
|
||||
* The terminal of the caller on this call's connection, or {@code null} when that caller carries
|
||||
* no terminal, which is only the unnamed primary. A named lead, an architect and a worker each
|
||||
* carry one.
|
||||
*/
|
||||
private static String callerTerminal(McpSyncServerExchange exchange) {
|
||||
Object v = exchange.transportContext().get(CALLER_TERMINAL);
|
||||
String s = v == null ? null : v.toString();
|
||||
@@ -876,7 +908,8 @@ public final class FleetMcp {
|
||||
* a BUSY interloper never claims a turn it did not win. {@code null} disables recording.
|
||||
*/
|
||||
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
|
||||
Long timeoutMs, Runnable onAccepted, Set<String> profiles) {
|
||||
Long timeoutMs, Runnable onAccepted, Set<String> profiles,
|
||||
String callerTerminal) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
@@ -886,7 +919,7 @@ public final class FleetMcp {
|
||||
}
|
||||
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
||||
try {
|
||||
return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout);
|
||||
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerTerminal), timeout);
|
||||
} catch (HerdrException e) {
|
||||
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
|
||||
}
|
||||
@@ -896,13 +929,15 @@ public final class FleetMcp {
|
||||
* {@code fleet_send} carrying a {@code turnId}: the primary's answer to a worker's
|
||||
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
|
||||
* it resumes the same turn — surfaced to the primary identically to a normal send.
|
||||
* {@code callerTerminal} must match the turn's recorded owner or this is refused.
|
||||
*/
|
||||
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs) {
|
||||
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs,
|
||||
String callerTerminal) {
|
||||
if (isBlank(turnId) || isBlank(content)) {
|
||||
return error("turnId and content are required to answer a worker's question");
|
||||
}
|
||||
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
||||
return formatReply(messages.answer(turnId, content, timeout), timeout);
|
||||
return formatReply(messages.answer(turnId, content, timeout, callerTerminal), timeout);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -951,6 +986,8 @@ public final class FleetMcp {
|
||||
+ "\" and content set to your answer; the worker resumes the same turn.");
|
||||
case STALE_TURN -> error("that question is no longer open — it timed out or was already "
|
||||
+ "answered (turnId stale)");
|
||||
case NOT_TURN_OWNER -> error("this turn belongs to a different delegation — only the caller "
|
||||
+ "that opened it may answer it");
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
|
||||
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
|
||||
// fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call,
|
||||
@@ -972,6 +1009,17 @@ public final class FleetMcp {
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
|
||||
Runnable onAccepted, Set<String> profiles) {
|
||||
return sendAsync(messages, sessionId, content, onAccepted, profiles, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
|
||||
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
|
||||
* result back to that same caller — see {@link MessageService#poll(String, String)}.
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
|
||||
Runnable onAccepted, Set<String> profiles,
|
||||
String creatorTerminal) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
@@ -979,7 +1027,7 @@ public final class FleetMcp {
|
||||
if (targetError != null) {
|
||||
return targetError;
|
||||
}
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted);
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
|
||||
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
|
||||
}
|
||||
|
||||
@@ -1158,9 +1206,24 @@ public final class FleetMcp {
|
||||
* exclusively of {@code ticket}/{@code target} — see {@link #pollAction}'s javadoc for why this
|
||||
* is a different inbox (this daemon's own {@link LeadChannel}) that authorizes differently
|
||||
* ({@link Authz.Action#COORD_READ}, primary-only) from either of the original two branches.
|
||||
*
|
||||
* <p>Does not check who owns {@code ticket} — see the overload that takes {@code
|
||||
* callerTerminal} for that. Callers that do not resolve a caller terminal (tests, or a surface
|
||||
* with no connection-based identity) use this one.
|
||||
*/
|
||||
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
|
||||
String target, String coordId) {
|
||||
return poll(messages, leadChannel, ticket, target, coordId, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
|
||||
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
|
||||
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
|
||||
* client-supplied value.
|
||||
*/
|
||||
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
|
||||
String target, String coordId, String callerTerminal) {
|
||||
if (!isBlank(coordId)) {
|
||||
return pollHeldPeerMail(leadChannel, coordId);
|
||||
}
|
||||
@@ -1174,7 +1237,7 @@ public final class FleetMcp {
|
||||
if (isBlank(ticket)) {
|
||||
return error("ticket (or target) is required");
|
||||
}
|
||||
MessageService.TaskView v = messages.poll(ticket);
|
||||
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
|
||||
if (v == null) {
|
||||
return error("unknown ticket: " + ticket + " (never issued, or expired)");
|
||||
}
|
||||
@@ -1281,16 +1344,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);
|
||||
}
|
||||
@@ -1778,7 +1845,7 @@ public final class FleetMcp {
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
|
||||
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
|
||||
Map.of(), false, coordination, false);
|
||||
Map.of(), false, coordination, false, true, true);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1788,7 +1855,7 @@ public final class FleetMcp {
|
||||
Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
|
||||
Map.of(), false, CoordinationSource.none(), false);
|
||||
Map.of(), false, CoordinationSource.none(), false, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1806,7 +1873,7 @@ public final class FleetMcp {
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
|
||||
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
|
||||
Map.of(), false, coordination, false);
|
||||
Map.of(), false, coordination, false, true, true);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1816,7 +1883,7 @@ public final class FleetMcp {
|
||||
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
|
||||
Map.of(), false, coordination, false);
|
||||
Map.of(), false, coordination, false, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1839,7 +1906,7 @@ public final class FleetMcp {
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
|
||||
Map.of(), false, coordination, false);
|
||||
Map.of(), false, coordination, false, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1856,15 +1923,23 @@ public final class FleetMcp {
|
||||
* {@code false} (fleetd #463: a forgotten argument fails closed, not
|
||||
* open), so a test that wants the {@code coordinator} row must pass
|
||||
* an explicit {@code true}
|
||||
* @param leadsVisible whether this caller may see the {@code leads} array (see
|
||||
* {@link #leadsVisibleTo}) — unlike {@code callerIsPrimary}, this is
|
||||
* also {@code true} for an architect or a collaborator, so it cannot
|
||||
* be derived from {@code callerIsPrimary} alone
|
||||
* @param membersVisible whether this caller may see the {@code members} array (see
|
||||
* {@link #membersVisibleTo}); also {@code true} for an architect, but
|
||||
* unlike {@code leadsVisible}, never for a collaborator
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
CoordinationSource coordination, boolean callerIsPrimary,
|
||||
boolean leadsVisible, boolean membersVisible) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
|
||||
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
|
||||
Map.of(), false, coordination, callerIsPrimary);
|
||||
Map.of(), false, coordination, callerIsPrimary, leadsVisible, membersVisible);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1881,6 +1956,10 @@ public final class FleetMcp {
|
||||
* {@link #collaboratorsVisibleTo}); every wrapper overload above
|
||||
* passes {@code false}, so a test that wants the row must call this
|
||||
* overload with an explicit {@code true}
|
||||
* @param leadsVisible whether this caller may see the {@code leads} array (see
|
||||
* {@link #leadsVisibleTo})
|
||||
* @param membersVisible whether this caller may see the {@code members} array (see
|
||||
* {@link #membersVisibleTo})
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
@@ -1890,28 +1969,41 @@ public final class FleetMcp {
|
||||
LeadConfigDirSource leadConfigDirs,
|
||||
Map<String, String> leads, String selfTerm,
|
||||
Map<String, String> collaborators, boolean collaboratorsVisible,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
CoordinationSource coordination, boolean callerIsPrimary,
|
||||
boolean leadsVisible, boolean membersVisible) {
|
||||
try {
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||
List<Map<String, Object>> leadRows = leads.entrySet().stream()
|
||||
.sorted(Map.Entry.comparingByValue())
|
||||
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
|
||||
leadConfigDirs))
|
||||
.toList();
|
||||
// Neither row's assembly (leadView/memberCapacityView probing herdr for live status)
|
||||
// runs unless at least one of them needs the live-agent lookup backing it.
|
||||
Map<String, Agent> live = (leadsVisible || membersVisible)
|
||||
? workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b))
|
||||
: Map.of();
|
||||
// fleetd #209: this is the caller-driven fleet_list read that actually reports
|
||||
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
|
||||
// resolving roster; the heartbeat/health/metrics timers stay on the plain sessions.roster().
|
||||
List<MemberSession> roster = sessions.rosterResolved();
|
||||
List<Map<String, Object>> out = roster.stream()
|
||||
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
|
||||
.toList();
|
||||
Set<String> profiles = new java.util.TreeSet<>(capacity.configuredProfiles().get());
|
||||
roster.stream().map(MemberSession::profile).forEach(profiles::add);
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
result.put("leads", leadRows); result.put("members", out);
|
||||
// READ is permission to enter this tool, not permission to receive every field it can
|
||||
// build -- gate BEFORE assembling each row, so the key is absent rather than
|
||||
// present-and-empty; a caller without either row gets its own facts from fleet_whoami.
|
||||
if (leadsVisible) {
|
||||
List<Map<String, Object>> leadRows = leads.entrySet().stream()
|
||||
.sorted(Map.Entry.comparingByValue())
|
||||
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
|
||||
leadConfigDirs))
|
||||
.toList();
|
||||
result.put("leads", leadRows);
|
||||
}
|
||||
if (membersVisible) {
|
||||
List<Map<String, Object>> out = roster.stream()
|
||||
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
|
||||
.toList();
|
||||
result.put("members", out);
|
||||
}
|
||||
result.put("healthCoverage", healthCoverage.value().get());
|
||||
result.put("loopHealth", Map.of(
|
||||
"statusPoller", loopHealth.statusPoller().get().name(),
|
||||
@@ -2445,7 +2537,14 @@ public final class FleetMcp {
|
||||
|
||||
private static McpSchema.Tool listTool() {
|
||||
return tool(FleetTool.LIST.wireName(),
|
||||
"List the whole fleet the bridge tracks, in two parts. 'leads' are your PEERS — other "
|
||||
"List the whole fleet the bridge tracks, in two parts. 'members' is visible to "
|
||||
+ "the primary and an architect only. 'leads' is visible to those two AND a "
|
||||
+ "collaborator — exactly the roles that may fleet_send to a lead, so a "
|
||||
+ "collaborator can learn a lead's sessionId before using the send it already "
|
||||
+ "holds. A worker holds READ to call this tool at all, but gets neither "
|
||||
+ "array, never an empty one; a worker reads its own session, "
|
||||
+ "profile, state, worktree, branch and owner from fleet_whoami instead. "
|
||||
+ "'leads' are your PEERS — other "
|
||||
+ "orchestrators, each with its sessionId (the address to fleet_send to), "
|
||||
+ "name, live status, and 'self': true on your own row; this is how you "
|
||||
+ "discover a peer lead without being told its address. 'members' are the "
|
||||
|
||||
@@ -87,7 +87,7 @@ public final class MessageService {
|
||||
/**
|
||||
* The worker paused mid-turn to ask the primary a question (CB-205); {@code text} is the
|
||||
* question and {@code turnId} correlates the answer. Not terminal — the primary answers with
|
||||
* {@link #answer(String, String, long)} and the turn resumes.
|
||||
* {@link #answer(String, String, long, String)} and the turn resumes.
|
||||
*/
|
||||
QUESTION,
|
||||
/** Timed out after the message was delivered — the worker is still working. */
|
||||
@@ -120,10 +120,16 @@ public final class MessageService {
|
||||
/** Another send to this session was in flight for the whole window. */
|
||||
BUSY,
|
||||
/**
|
||||
* An answer ({@link #answer(String, String, long)}) referenced a {@code turnId} that is no
|
||||
* longer open — the worker's {@code fleet_ask} already timed out or was answered.
|
||||
* An answer ({@link #answer(String, String, long, String)}) referenced a {@code turnId}
|
||||
* that is no longer open — the worker's {@code fleet_ask} already timed out or was answered.
|
||||
*/
|
||||
STALE_TURN
|
||||
STALE_TURN,
|
||||
/**
|
||||
* An answer ({@link #answer(String, String, long, String)}) named a {@code turnId} that is
|
||||
* still open, but the answering caller is not the caller whose accepted delegation opened
|
||||
* it. Distinct from {@link #STALE_TURN} so a refusal is never reported as a lapsed turn.
|
||||
*/
|
||||
NOT_TURN_OWNER
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -133,7 +139,7 @@ public final class MessageService {
|
||||
* {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION},
|
||||
* else {@code null}
|
||||
* @param turnId correlation id for a {@link Outcome#QUESTION} (answered via
|
||||
* {@link #answer(String, String, long)}), else {@code null}
|
||||
* {@link #answer(String, String, long, String)}), else {@code null}
|
||||
*/
|
||||
public record Reply(Outcome outcome, String text, String turnId) {
|
||||
/** A reply with no correlation id (the common terminal outcomes). */
|
||||
@@ -275,10 +281,19 @@ public final class MessageService {
|
||||
* {@link #abandon}) can never match again regardless of this flag's value.
|
||||
*/
|
||||
private volatile boolean askTimedOut;
|
||||
/**
|
||||
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
|
||||
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
|
||||
* created through an overload that does not record one. {@link #poll(String, String)}
|
||||
* compares a polling caller's own terminal against this field before handing back the
|
||||
* ticket's state.
|
||||
*/
|
||||
private final String creatorTerminal;
|
||||
|
||||
private Task(String ticket, String target, LongSupplier nowNanos) {
|
||||
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
|
||||
this.ticket = ticket;
|
||||
this.target = target;
|
||||
this.creatorTerminal = creatorTerminal;
|
||||
this.createdNanos = nowNanos.getAsLong();
|
||||
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
|
||||
}
|
||||
@@ -656,7 +671,7 @@ public final class MessageService {
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, TIMED_OUT_UNCONFIRMED, BUSY -> "timeout";
|
||||
case WORKER_FAILED -> "failed";
|
||||
case BACKEND_EXHAUSTED -> "backend_exhausted";
|
||||
case STALE_TURN, QUESTION -> null; // not a completed delegation
|
||||
case STALE_TURN, QUESTION, NOT_TURN_OWNER -> null; // not a completed delegation
|
||||
};
|
||||
}
|
||||
|
||||
@@ -915,14 +930,17 @@ public final class MessageService {
|
||||
|
||||
/**
|
||||
* Deliver {@code content} to {@code target} (a herdr {@code terminal_id}) and block until the
|
||||
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses.
|
||||
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. {@code callerTerminal}
|
||||
* is the terminal of the caller making this call — {@code null} for the unnamed primary — and is
|
||||
* recorded as the turn's owner, the only caller {@link #answer(String, String, long, String)} will
|
||||
* later accept an answer from if the worker pauses mid-turn to ask.
|
||||
*/
|
||||
public Reply send(String target, String content, long timeoutMillis) {
|
||||
return send(target, content, timeoutMillis, null);
|
||||
public Reply send(String target, String content, long timeoutMillis, String callerTerminal) {
|
||||
return send(target, content, timeoutMillis, null, callerTerminal);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #send(String, String, long)}, but with an accepted-delivery hook.
|
||||
* As {@link #send(String, String, long, String)}, but with an accepted-delivery hook.
|
||||
*
|
||||
* <p>{@code onAccepted} is invoked exactly once, once this send has won {@code target}'s send
|
||||
* lock and so become the <em>accepted target turn</em> — it runs <em>before</em> delivery is
|
||||
@@ -933,12 +951,13 @@ public final class MessageService {
|
||||
* acceptance means a concurrent sender that times out {@code BUSY} can never steal ownership it
|
||||
* never earned. {@code null} disables the hook.
|
||||
*/
|
||||
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted) {
|
||||
return send(target, content, timeoutMillis, onAccepted, null);
|
||||
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerTerminal) {
|
||||
return send(target, content, timeoutMillis, onAccepted, null, callerTerminal);
|
||||
}
|
||||
|
||||
/** Run a send, optionally stopping an async task that teardown already failed before acceptance. */
|
||||
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task) {
|
||||
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task,
|
||||
String callerTerminal) {
|
||||
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
|
||||
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
|
||||
|
||||
@@ -958,7 +977,7 @@ public final class MessageService {
|
||||
// open race). Opening first also means a throwing onAccepted (fired before enqueue) or an
|
||||
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
|
||||
// failed send leaves no stale waiter behind.
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target, Rendezvous.Owner.of(callerTerminal));
|
||||
// CB-640: this send now owns target's delivery, so any earlier stranded-reply or
|
||||
// still-queued fact no longer describes the live state — clear both rather than let
|
||||
// them outlive the send that supersedes them.
|
||||
@@ -1137,19 +1156,30 @@ public final class MessageService {
|
||||
* mid-turn (already picked up), so the answer flows back through its own open {@code fleet_ask}
|
||||
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
|
||||
* is unblocked so a reply that lands the instant it resumes is not lost.
|
||||
*
|
||||
* <p>{@code callerTerminal} is the terminal of the caller making this call — {@code null} for
|
||||
* the unnamed primary. It is checked against the turn's recorded owner (the caller whose
|
||||
* accepted delegation opened it, see {@link #send(String, String, long, String)} and
|
||||
* {@link #sendAsync(String, String, Runnable, String)}) before anything else runs: a mismatch,
|
||||
* including a turn with no owner on record at all, returns {@link Outcome#NOT_TURN_OWNER}
|
||||
* without touching the rendezvous, the session lock, or any async task bookkeeping.
|
||||
*/
|
||||
public Reply answer(String turnId, String content, long timeoutMillis) {
|
||||
public Reply answer(String turnId, String content, long timeoutMillis, String callerTerminal) {
|
||||
String workerSession = rendezvous.askSession(turnId);
|
||||
if (workerSession == null) {
|
||||
return new Reply(Outcome.STALE_TURN, null); // the ask lapsed (timed out or already answered)
|
||||
}
|
||||
Rendezvous.Owner owner = rendezvous.askOwner(turnId);
|
||||
if (!Rendezvous.Owner.permits(owner, callerTerminal)) {
|
||||
return new Reply(Outcome.NOT_TURN_OWNER, null);
|
||||
}
|
||||
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
|
||||
ReentrantLock lock = sessionLocks.computeIfAbsent(workerSession, _ -> new ReentrantLock());
|
||||
if (!tryLock(lock, remainingMillis(deadlineNanos))) {
|
||||
return new Reply(Outcome.BUSY, null);
|
||||
}
|
||||
try {
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession, owner);
|
||||
// fleetd #575: this try used to open below, AFTER the Task lookup/registration and the
|
||||
// STALE_TURN early return that follows it — so that return was covered only by a
|
||||
// hand-rolled copy of the finally's own cleanup pair, not the finally itself. Widening the
|
||||
@@ -1279,8 +1309,20 @@ public final class MessageService {
|
||||
* @return the ticket to poll for the eventual result
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted) {
|
||||
return sendAsync(target, content, onAccepted, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
|
||||
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
|
||||
* differs from this one; {@code null} records no owner (a caller with no terminal — the
|
||||
* unnamed primary — is always allowed to poll the result regardless).
|
||||
*
|
||||
* @return the ticket to poll for the eventual result
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
|
||||
String ticket = "task-" + ticketSeq.incrementAndGet();
|
||||
Task task = new Task(ticket, target, nowNanos);
|
||||
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
|
||||
tasks.put(ticket, task);
|
||||
if (pushLoop != null) {
|
||||
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
||||
@@ -1311,7 +1353,7 @@ public final class MessageService {
|
||||
}
|
||||
asyncExecutor.submit(() -> {
|
||||
try {
|
||||
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
|
||||
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorTerminal);
|
||||
if (result.outcome() == Outcome.QUESTION) {
|
||||
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
|
||||
// just after resolveQuestion wakes this thread.
|
||||
@@ -1343,15 +1385,31 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket;
|
||||
* otherwise a {@link Phase#PENDING} view (with the live worker status as detail), a
|
||||
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
|
||||
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
|
||||
* checked, so this overload must only be used where the caller's identity is otherwise
|
||||
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
|
||||
*/
|
||||
public TaskView poll(String ticket) {
|
||||
return poll(ticket, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
|
||||
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
|
||||
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
|
||||
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
|
||||
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
|
||||
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
|
||||
*/
|
||||
public TaskView poll(String ticket, String callerTerminal) {
|
||||
Task task = tasks.get(ticket);
|
||||
if (task == null) {
|
||||
return null;
|
||||
}
|
||||
if (!ownsTicket(task, callerTerminal)) {
|
||||
return new TaskView(ticket, Phase.FAILED, null, null,
|
||||
"forbidden: this ticket was created by a different session", null);
|
||||
}
|
||||
CompletableFuture<Reply> f = task.future;
|
||||
if (!f.isDone()) {
|
||||
Reply question = task.question;
|
||||
@@ -1387,6 +1445,17 @@ public final class MessageService {
|
||||
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
|
||||
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
|
||||
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
|
||||
* equal the terminal recorded on the task; a task with no recorded terminal matches no
|
||||
* terminal-bearing caller.
|
||||
*/
|
||||
private static boolean ownsTicket(Task task, String callerTerminal) {
|
||||
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
|
||||
}
|
||||
|
||||
/**
|
||||
* Test seam only — carries no production behaviour, and nothing in this class calls it;
|
||||
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
|
||||
@@ -1706,16 +1775,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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -60,7 +60,7 @@ public final class Rendezvous {
|
||||
}
|
||||
|
||||
/** A worker's open mid-turn question: the worker session it belongs to and the answer future. */
|
||||
private record AskWaiter(String session, CompletableFuture<String> answer) {
|
||||
private record AskWaiter(String session, CompletableFuture<String> answer, Owner owner) {
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -70,7 +70,33 @@ public final class Rendezvous {
|
||||
public record AskTicket(String turnId, CompletableFuture<String> answer, boolean fresh) {
|
||||
}
|
||||
|
||||
private final ConcurrentHashMap<String, CompletableFuture<Resolution>> waiters = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* The caller whose accepted delegation opened a turn — the only caller allowed to answer it.
|
||||
* A {@code null} terminal means the unnamed primary, an authenticated caller with no pane.
|
||||
*/
|
||||
public record Owner(String terminal) {
|
||||
public static final Owner UNNAMED_PRIMARY = new Owner(null);
|
||||
|
||||
public static Owner of(String terminal) {
|
||||
return terminal == null ? UNNAMED_PRIMARY : new Owner(terminal);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code callerTerminal} matches {@code owner}. A {@code null} owner means no
|
||||
* owner was ever recorded, and that state matches no caller, not even one whose own
|
||||
* terminal is {@code null} — "no record" and "recorded as the unnamed primary" are
|
||||
* different states.
|
||||
*/
|
||||
public static boolean permits(Owner owner, String callerTerminal) {
|
||||
return owner != null && java.util.Objects.equals(owner.terminal(), callerTerminal);
|
||||
}
|
||||
}
|
||||
|
||||
/** A registered forward waiter together with the owner its delegation was opened under. */
|
||||
private record ForwardWaiter(Owner owner, CompletableFuture<Resolution> future) {
|
||||
}
|
||||
|
||||
private final ConcurrentHashMap<String, ForwardWaiter> waiters = new ConcurrentHashMap<>();
|
||||
|
||||
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */
|
||||
private final ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
|
||||
@@ -89,8 +115,17 @@ public final class Rendezvous {
|
||||
* code a double open is impossible; this is a tripwire for the day that no longer holds.
|
||||
*/
|
||||
public CompletableFuture<Resolution> open(String session) {
|
||||
return open(session, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Same as {@link #open(String)}, additionally recording {@code owner} as the caller whose
|
||||
* delegation opened this waiter. A {@code null} owner records no owner at all — the
|
||||
* fail-closed default {@link Owner#permits} refuses to everyone.
|
||||
*/
|
||||
public CompletableFuture<Resolution> open(String session, Owner owner) {
|
||||
CompletableFuture<Resolution> waiter = new CompletableFuture<>();
|
||||
CompletableFuture<Resolution> existing = waiters.putIfAbsent(session, waiter);
|
||||
ForwardWaiter existing = waiters.putIfAbsent(session, new ForwardWaiter(owner, waiter));
|
||||
if (existing != null) {
|
||||
throw new IllegalStateException(
|
||||
"rendezvous double-open for session " + session + " — a waiter is already registered");
|
||||
@@ -105,7 +140,13 @@ public final class Rendezvous {
|
||||
* successful {@code open} after a finished turn requires this close to have happened first).
|
||||
*/
|
||||
public void close(String session, CompletableFuture<Resolution> waiter) {
|
||||
waiters.remove(session, waiter);
|
||||
waiters.computeIfPresent(session, (s, w) -> w.future() == waiter ? null : w);
|
||||
}
|
||||
|
||||
/** The owner recorded for {@code session}'s open waiter, or {@code null} if none is open. */
|
||||
public Owner ownerOf(String session) {
|
||||
ForwardWaiter w = waiters.get(session);
|
||||
return w == null ? null : w.owner();
|
||||
}
|
||||
|
||||
/** Whether a send is currently awaiting a resolution for {@code session}. */
|
||||
@@ -119,7 +160,8 @@ public final class Rendezvous {
|
||||
* send (see the CB-116 note above) rather than whichever send happens to be waiting when they fire.
|
||||
*/
|
||||
public CompletableFuture<Resolution> currentWaiter(String session) {
|
||||
return waiters.get(session);
|
||||
ForwardWaiter w = waiters.get(session);
|
||||
return w == null ? null : w.future();
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -147,7 +189,7 @@ public final class Rendezvous {
|
||||
String turnId = openAsksBySession.computeIfAbsent(session, _ -> {
|
||||
String newTurnId = session + "#" + askSeq.incrementAndGet();
|
||||
CompletableFuture<String> answer = new CompletableFuture<>();
|
||||
AskWaiter waiter = new AskWaiter(session, answer);
|
||||
AskWaiter waiter = new AskWaiter(session, answer, ownerOf(session));
|
||||
asks.put(newTurnId, waiter);
|
||||
minted[0] = waiter;
|
||||
return newTurnId;
|
||||
@@ -183,6 +225,16 @@ public final class Rendezvous {
|
||||
return w == null ? null : w.session();
|
||||
}
|
||||
|
||||
/**
|
||||
* The owner recorded for {@code turnId} when its ask turn was freshly opened — the caller
|
||||
* whose delegation {@link #answerAsk} must match. {@code null} if {@code turnId} is unknown or
|
||||
* lapsed, or if the ask opened with no forward waiter owner on record.
|
||||
*/
|
||||
public Owner askOwner(String turnId) {
|
||||
AskWaiter w = asks.get(turnId);
|
||||
return w == null ? null : w.owner();
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a worker's blocked {@code fleet_ask} with the primary's {@code answer}, unblocking it
|
||||
* to resume its turn.
|
||||
@@ -245,7 +297,7 @@ public final class Rendezvous {
|
||||
}
|
||||
|
||||
private boolean complete(String session, Resolution resolution) {
|
||||
CompletableFuture<Resolution> waiter = waiters.get(session);
|
||||
return waiter != null && waiter.complete(resolution);
|
||||
ForwardWaiter waiter = waiters.get(session);
|
||||
return waiter != null && waiter.future().complete(resolution);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -663,6 +663,8 @@ public final class FleetApp {
|
||||
if (!allow(ctx, routeAction("POST /sessions/{id}/message"), id)) {
|
||||
return;
|
||||
}
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
String callerTerminal = caller == null ? null : caller.terminal();
|
||||
JsonNode body;
|
||||
try {
|
||||
body = mapper.readTree(ctx.body());
|
||||
@@ -688,19 +690,19 @@ public final class FleetApp {
|
||||
|
||||
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
|
||||
if (turnId != null && !turnId.isBlank()) {
|
||||
writeReply(ctx, id, messages.answer(turnId, content, timeout), timeout);
|
||||
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerTerminal), timeout);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!wait) {
|
||||
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
|
||||
String ticket = messages.sendAsync(id, content);
|
||||
String ticket = messages.sendAsync(id, content, null, callerTerminal);
|
||||
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
writeReply(ctx, id, messages.send(id, content, timeout), timeout);
|
||||
writeReply(ctx, id, messages.send(id, content, timeout, callerTerminal), timeout);
|
||||
} catch (HerdrException e) {
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
@@ -719,6 +721,10 @@ public final class FleetApp {
|
||||
case STALE_TURN -> ctx.status(409).json(Map.of(
|
||||
"sessionId", id, "error", "stale_turn",
|
||||
"detail", "that question is no longer open (timed out or already answered)"));
|
||||
case NOT_TURN_OWNER -> ctx.status(403).json(Map.of(
|
||||
"sessionId", id, "error", "not_turn_owner",
|
||||
"detail", "this turn belongs to a different delegation — only the caller that "
|
||||
+ "opened it may answer it"));
|
||||
case REPLIED, COMPLETED_UNREPLIED -> {
|
||||
// replySource distinguishes a structured fleet_reply from the CB-106 completion
|
||||
// fallback (a scrape of the worker's transcript when it finished without replying).
|
||||
@@ -733,9 +739,10 @@ public final class FleetApp {
|
||||
// silent fall-through. That is exactly the bug this ticket exists to fix:
|
||||
// `default -> "done"` used to sit here and would have told a REST caller the
|
||||
// delegation completed for TIMED_OUT_UNCONFIRMED, the one outcome where delivery
|
||||
// is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION and STALE_TURN can never
|
||||
// actually reach this inner switch — the outer switch above always dispatches
|
||||
// them first — but they still need an arm to keep this switch exhaustive.
|
||||
// is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN and
|
||||
// NOT_TURN_OWNER can never actually reach this inner switch — the outer switch
|
||||
// above always dispatches them first — but they still need an arm to keep this
|
||||
// switch exhaustive.
|
||||
"status", switch (reply.outcome()) {
|
||||
case TIMED_OUT_WORKING -> "working";
|
||||
case TIMED_OUT_QUEUED -> "queued";
|
||||
@@ -745,7 +752,7 @@ public final class FleetApp {
|
||||
case BUSY -> "busy";
|
||||
case WORKER_FAILED -> "failed";
|
||||
case BACKEND_EXHAUSTED -> "backend_exhausted";
|
||||
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable
|
||||
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN, NOT_TURN_OWNER -> "done"; // unreachable
|
||||
},
|
||||
"detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED
|
||||
? "no reply within " + timeout + "ms; delivery is unconfirmed — the "
|
||||
@@ -874,10 +881,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 +903,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;
|
||||
|
||||
@@ -26,6 +26,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
/**
|
||||
* Authoritative in-daemon registry of the worker sessions this {@code fleetd} process spawned.
|
||||
@@ -57,6 +58,18 @@ public final class SessionManager implements TurnListener {
|
||||
* Populated on every spawn path, removed on {@link #release}.
|
||||
*/
|
||||
private final ConcurrentHashMap<String /*paneId*/, PeerHandle> handles = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* fleetd #702: a pane mid-teardown, keyed by paneId, held from just before its registry entry
|
||||
* is removed until {@link #releaseRemoved} finishes. {@link #spawnedMemberRole} consults this
|
||||
* alongside the registry, so a caller resolving the pane's terminal during that window still
|
||||
* sees a live member and never falls through to a tab map.
|
||||
*
|
||||
* <p>Depth-counted rather than a plain set: two threads can be tearing down the same pane at
|
||||
* once (the CAS in {@link #releaseIfCurrent} exists for exactly that race), and with a set the
|
||||
* loser's {@code finally} would unmark the pane while the winner is still mid-teardown,
|
||||
* reopening the window this exists to close.
|
||||
*/
|
||||
private final ConcurrentHashMap<String /*paneId*/, Releasing> releasing = new ConcurrentHashMap<>();
|
||||
private final MemberPresence presence;
|
||||
private final SecureRandom nonceRandom = new SecureRandom();
|
||||
private final AtomicLong nonceSeq = new AtomicLong();
|
||||
@@ -303,9 +316,11 @@ public final class SessionManager implements TurnListener {
|
||||
* is a logged path an operator can reclaim, the cost of a deleted one is unrecoverable work.
|
||||
*/
|
||||
private MemberSession release(String paneId, ReleaseCause cause) {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
|
||||
return removed;
|
||||
return releaseWindow(paneId, registry.get(paneId), () -> {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
|
||||
return removed;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -314,16 +329,85 @@ public final class SessionManager implements TurnListener {
|
||||
* DONE record from stopping a worker that delivery has made BUSY.
|
||||
*/
|
||||
private boolean releaseIfCurrent(MemberSession expected, ReleaseCause cause) {
|
||||
if (!registry.remove(expected.paneId(), expected)) {
|
||||
// A lifecycle transition replaced the record between the caller's check and this remove.
|
||||
// Log it: this race is by definition unobservable otherwise, and a reaper that silently
|
||||
// declines to reap is the hardest kind of behaviour to diagnose after the fact.
|
||||
log.debug("skipping reap of pane={}: its registry record changed after the idle check "
|
||||
+ "(most likely a delivery made it BUSY)", expected.paneId());
|
||||
return false;
|
||||
return releaseWindow(expected.paneId(), expected, () -> {
|
||||
if (!registry.remove(expected.paneId(), expected)) {
|
||||
// A lifecycle transition replaced the record between the caller's check and this
|
||||
// remove. Log it: this race is by definition unobservable otherwise, and a reaper
|
||||
// that silently declines to reap is the hardest kind of behaviour to diagnose
|
||||
// after the fact.
|
||||
log.debug("skipping reap of pane={}: its registry record changed after the idle "
|
||||
+ "check (most likely a delivery made it BUSY)", expected.paneId());
|
||||
return false;
|
||||
}
|
||||
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
|
||||
return true;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #702: mark {@code paneId} as mid-teardown — using {@code known}'s terminal/role when
|
||||
* it is available — for the whole of {@code teardown}, which removes the registry entry and
|
||||
* then runs {@link #releaseRemoved}. Shared by both registry-removal sites ({@link #release}'s
|
||||
* unconditional remove and {@link #releaseIfCurrent}'s CAS remove) so neither can leave the
|
||||
* other's window unmarked.
|
||||
*
|
||||
* <p>The mark is written before {@code teardown} runs — so it covers the removal itself, not
|
||||
* only what comes after it — and cleared in a {@code finally}, so an unchecked throw out of
|
||||
* {@code teardown} (including one from {@link PeerLauncher#stop}, which declares nothing) can
|
||||
* never leave the pane marked for the rest of the daemon's life.
|
||||
*/
|
||||
private <T> T releaseWindow(String paneId, MemberSession known, Supplier<T> teardown) {
|
||||
releasing.compute(paneId, (_, prior) -> Releasing.enter(prior, known));
|
||||
try {
|
||||
return teardown.get();
|
||||
} finally {
|
||||
releasing.compute(paneId, (_, prior) -> prior == null ? null : prior.leave());
|
||||
}
|
||||
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Depth count plus the terminal/role a mid-teardown pane belongs to, for
|
||||
* {@link #spawnedMemberRole}. The terminal/role come from whichever call into
|
||||
* {@link #releaseWindow} first knew them: a call that finds the registry entry already gone
|
||||
* passes a {@code null} session, and must not blank out what the first call recorded.
|
||||
*/
|
||||
record Releasing(int depth, String terminalId, MemberRole role) {
|
||||
static Releasing enter(Releasing prior, MemberSession known) {
|
||||
int depth = (prior == null ? 0 : prior.depth()) + 1;
|
||||
String terminalId = known != null ? known.terminalId() : prior == null ? null : prior.terminalId();
|
||||
MemberRole role = known != null ? known.role() : prior == null ? null : prior.role();
|
||||
return new Releasing(depth, terminalId, role);
|
||||
}
|
||||
|
||||
Releasing leave() {
|
||||
return depth <= 1 ? null : new Releasing(depth - 1, terminalId, role);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The role of the live spawned member occupying {@code terminal} — whether it is currently in
|
||||
* the registry, or mid-teardown between {@link #release} removing its registry entry and
|
||||
* {@link #releaseRemoved} actually stopping its pane (fleetd #702). {@code null} for a terminal
|
||||
* that is neither: this method is the one reader a caller resolver consults before any tab
|
||||
* map, so a live or releasing member's identity never falls back to a tab label.
|
||||
*
|
||||
* <p>Checks the registry directly via {@link #findByTerminal} rather than {@link #roster()},
|
||||
* so this hot-path lookup (consulted on every resolve) never pays for a list copy or a stream.
|
||||
*/
|
||||
public MemberRole spawnedMemberRole(String terminal) {
|
||||
MemberSession session = findByTerminal(terminal);
|
||||
if (session != null) {
|
||||
return session.role();
|
||||
}
|
||||
if (terminal == null) {
|
||||
return null;
|
||||
}
|
||||
for (Releasing r : releasing.values()) {
|
||||
if (terminal.equals(r.terminalId())) {
|
||||
return r.role();
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private void releaseRemoved(String paneId, MemberSession removed, PeerHandle removedHandle,
|
||||
|
||||
@@ -30,6 +30,22 @@ class AuthzTest {
|
||||
assertTrue(Authz.isUnauthenticated(null));
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code MessageService.answer}'s turn-ownership check treats a caller with no terminal as
|
||||
* matching a turn recorded for the unnamed primary. That rule only stays safe because an
|
||||
* unauthenticated caller — whose terminal is also {@code null} — never reaches {@code answer}
|
||||
* at all: {@link #anonymousIsAuthorizedForNothing} already covers every action including
|
||||
* {@code ANSWER}, but this test names the exact coupling so a future change to either side
|
||||
* cannot drift without turning this test red.
|
||||
*/
|
||||
@Test
|
||||
void anAnonymousCallerIsRefusedAnswerSoItCanNeverBeMistakenForTheUnnamedPrimary() {
|
||||
assertFalse(Authz.permits(ANON, ANSWER, null),
|
||||
"an anonymous caller, whose terminal is also null, must never reach answer() — the "
|
||||
+ "turn-ownership check's null-terminal match for the unnamed primary owner "
|
||||
+ "relies on this gate refusing it first");
|
||||
}
|
||||
|
||||
@Test
|
||||
void orchestrationBelongsToThePrimaryAlone() {
|
||||
for (Authz.Action a : new Authz.Action[]{SPAWN, STOP, SEND, DRAIN}) {
|
||||
|
||||
@@ -1,13 +1,25 @@
|
||||
package dev.ltms.fleet.auth;
|
||||
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.mcp.ConnectionIdentity;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.WorktreeRequest;
|
||||
import dev.ltms.fleet.session.Worktrees;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -460,7 +472,6 @@ class CallerResolverTest {
|
||||
assertTrue(members.bind("architect:lead-designer", "term_a")); // the later lifecycle binds the slot
|
||||
|
||||
assertEquals(Role.ARCHITECT, r.resolve("127.0.0.1", 42, null).role());
|
||||
assertEquals("architect:lead-designer", r.members().get("term_a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -542,6 +553,219 @@ class CallerResolverTest {
|
||||
assertEquals("term_a", p.terminal());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #702: between {@link SessionManager#release} removing the registry entry and the
|
||||
* pane actually stopping, a resolve for that terminal must still see the live member and
|
||||
* never fall through to a tab map — wired through the real {@link SessionManager}, not a
|
||||
* hand-rolled stand-in for its {@code spawnedMemberRole}.
|
||||
*
|
||||
* <p>Reuses the ticket's own test idea: an injected {@code hasUncommitted} resolves the
|
||||
* releasing terminal from inside {@code release}'s window — a real call landing inside the
|
||||
* window, so no sleep and no race.
|
||||
*
|
||||
* <p>The property asserted is "no tab map is consulted", never "the same role is returned".
|
||||
* The mandatory control is the lead tab map: it names this exact terminal, and the same
|
||||
* resolve taken <em>outside</em> the window (before release runs) must still return the lead
|
||||
* role — without that control, the in-window assertion would also pass on an empty map and
|
||||
* prove nothing.
|
||||
*/
|
||||
@Test
|
||||
void aPaneMidTeardownResolvesAsItsOwnRoleConsultingNoTabMap() {
|
||||
FakeHerdr herdr = new FakeHerdr().pinNextStarts(1, "term_a", "w2:p7");
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
|
||||
AtomicReference<Principal> duringWindow = new AtomicReference<>();
|
||||
AtomicReference<CallerResolver> resolverRef = new AtomicReference<>();
|
||||
Worktrees worktrees = new Worktrees() {
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
return "/wt/" + branch.replace('/', '_');
|
||||
}
|
||||
|
||||
@Override
|
||||
public void remove(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
// Runs from INSIDE release()'s git-status shell-out: the registry entry is
|
||||
// already gone, but the pane has not stopped yet.
|
||||
duringWindow.set(resolverRef.get().resolve("127.0.0.1", 42, null));
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String repoRoot(String cwd) {
|
||||
return "/repo";
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<String> snapshot(String worktreePath, String branch, String message) {
|
||||
return Optional.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WipRefStats wipRefs(String repoRoot) {
|
||||
return new WipRefStats(0, 0L);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
}
|
||||
};
|
||||
SessionManager sessions = new SessionManager(workers, worktrees);
|
||||
|
||||
// The lead tab map names "term_a" before anything is ever spawned onto it — the control
|
||||
// this test needs. Built up front so the SAME resolver answers every resolve() call below.
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> FakeHerdr.WORKER_PID);
|
||||
CallerResolver resolver = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
() -> Map.of("term_a", "the-lead"), null, sessions::spawnedMemberRole, Map::of);
|
||||
resolverRef.set(resolver);
|
||||
|
||||
Principal before = resolver.resolve("127.0.0.1", 42, null);
|
||||
assertEquals(Role.PRIMARY, before.role(),
|
||||
"control: with no live or releasing member on this terminal, the lead tab map must "
|
||||
+ "win — this is what proves the in-window assertion below is not passing on "
|
||||
+ "an empty map");
|
||||
assertEquals("the-lead", before.name());
|
||||
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-702", null));
|
||||
assertEquals("term_a", s.terminalId(), "sanity: the spawn resolved to the pinned pane");
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertNotNull(duringWindow.get(), "the dirty check must have run and captured a resolve");
|
||||
assertEquals(Role.WORKER, duringWindow.get().role(),
|
||||
"inside the window the pane must resolve as its own live-member role, consulting no "
|
||||
+ "tab map — a lead tab naming the same terminal must not win");
|
||||
assertEquals("term_a", duringWindow.get().terminal());
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link SessionManager#releaseRemoved} unbinds the architect slot before the git-status
|
||||
* shell-out that opens the teardown window, so a resolve landing inside that window must see
|
||||
* the slot already unbound and resolve {@link Role#WORKER} — never {@link Role#ARCHITECT},
|
||||
* and never by checking role equality against the live session, which would hold even if the
|
||||
* unbind ran too late.
|
||||
*
|
||||
* <p>The control is the same resolve taken outside the window, while the slot is still bound,
|
||||
* which must return {@link Role#ARCHITECT} — without it this test would also pass against a
|
||||
* slot that was never bound, and prove nothing.
|
||||
*/
|
||||
@Test
|
||||
void aReleasingArchitectIsDemotedToWorkerInsideTheTeardownWindow() {
|
||||
FakeHerdr herdr = new FakeHerdr().pinNextStarts(1, "term_a", "w2:p7");
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
|
||||
AtomicReference<Principal> duringWindow = new AtomicReference<>();
|
||||
AtomicReference<CallerResolver> resolverRef = new AtomicReference<>();
|
||||
Worktrees worktrees = new Worktrees() {
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
return "/wt/" + branch.replace('/', '_');
|
||||
}
|
||||
|
||||
@Override
|
||||
public void remove(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
// Runs from INSIDE release()'s git-status shell-out: the architect slot is already
|
||||
// unbound by this point, but the pane has not stopped yet.
|
||||
duringWindow.set(resolverRef.get().resolve("127.0.0.1", 42, null));
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String repoRoot(String cwd) {
|
||||
return "/repo";
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<String> snapshot(String worktreePath, String branch, String message) {
|
||||
return Optional.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WipRefStats wipRefs(String repoRoot) {
|
||||
return new WipRefStats(0, 0L);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
}
|
||||
};
|
||||
SessionManager sessions = new SessionManager(workers, worktrees);
|
||||
|
||||
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(Map.of(),
|
||||
Map.of("lead-designer", new FleetConfig.Slot("ltms-local")), Map.of(), Map.of(), null));
|
||||
sessions.setMemberLifecycle(members);
|
||||
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> FakeHerdr.WORKER_PID);
|
||||
CallerResolver resolver = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, members, sessions::spawnedMemberRole, Map::of);
|
||||
resolverRef.set(resolver);
|
||||
|
||||
MemberSession s = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-702d", null));
|
||||
assertEquals("term_a", s.terminalId(), "sanity: the spawn resolved to the pinned pane");
|
||||
assertEquals(MemberRole.ARCHITECT, s.role(), "sanity: the slot bind succeeded");
|
||||
|
||||
Principal before = resolver.resolve("127.0.0.1", 42, null);
|
||||
assertEquals(Role.ARCHITECT, before.role(),
|
||||
"control: with the slot still bound, the pane must resolve as an architect — this "
|
||||
+ "is what proves the in-window assertion below is not passing against a "
|
||||
+ "slot that was never bound");
|
||||
assertEquals("lead-designer", before.name());
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertNotNull(duringWindow.get(), "the dirty check must have run and captured a resolve");
|
||||
assertEquals(Role.WORKER, duringWindow.get().role(),
|
||||
"inside the window the architect slot is already unbound, so the result must be "
|
||||
+ "WORKER — asserting role-equality with the live session here would tempt "
|
||||
+ "moving the unbind earlier or later, which would be wrong either way");
|
||||
assertEquals("term_a", duringWindow.get().terminal());
|
||||
}
|
||||
|
||||
/** A spawned architect in the roster resolves ARCHITECT, carrying its bound slot's name. */
|
||||
@Test
|
||||
void aSpawnedArchitectInTheRosterResolvesArchitectWithItsSlotName() {
|
||||
|
||||
@@ -355,10 +355,10 @@ class FleetMcpAuthzTest {
|
||||
+ "the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
Pattern trailingArg = Pattern.compile(
|
||||
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*\\)\\s*;",
|
||||
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*[,)]",
|
||||
Pattern.DOTALL);
|
||||
Matcher m = trailingArg.matcher(handlerBlock);
|
||||
assertTrue(m.find(), "could not locate listFleet(...)'s trailing boolean argument in the "
|
||||
assertTrue(m.find(), "could not locate listFleet(...)'s coordinator boolean argument in the "
|
||||
+ "listHandler block -- the call shape changed, update this test's anchor: " + handlerBlock);
|
||||
String trailing = m.group(1);
|
||||
assertEquals("coordinatorVisibleTo(principal(exchange))", trailing,
|
||||
@@ -418,6 +418,145 @@ class FleetMcpAuthzTest {
|
||||
+ "calling, not pass a literal boolean -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
// --- who may see fleet_list's leads and members arrays ---------------------------------------
|
||||
|
||||
/**
|
||||
* {@link FleetMcp#leadsVisibleTo} is the whole policy decision for {@code fleet_list}'s
|
||||
* {@code leads} array: visible to exactly the roles that may {@code SEND} to a lead -- the
|
||||
* primary, an architect, and a collaborator -- never a worker, which holds {@code READ} but
|
||||
* can never {@code SEND} at all, and never an anonymous caller.
|
||||
*/
|
||||
@Test
|
||||
void primaryArchitectAndCollaboratorMaySeeTheLeadsArray() {
|
||||
assertTrue(FleetMcp.leadsVisibleTo(PRIMARY), "the primary must see the leads array");
|
||||
assertTrue(FleetMcp.leadsVisibleTo(ARCH_DESIGN), "an architect must see the leads array");
|
||||
assertTrue(FleetMcp.leadsVisibleTo(COLLABORATOR),
|
||||
"a collaborator may SEND to a lead, so it must see the leads array to learn where");
|
||||
assertFalse(FleetMcp.leadsVisibleTo(WORKER_A),
|
||||
"a worker holds READ but can never SEND, so it must not see the leads array");
|
||||
assertFalse(FleetMcp.leadsVisibleTo(ANON), "authenticated as nothing must not see it either");
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #primaryArchitectAndCollaboratorMaySeeTheLeadsArray}, for the {@code members}
|
||||
* array -- but a collaborator may {@code SEND} only to a lead or another collaborator, never
|
||||
* to a spawned member, so it must not see this one.
|
||||
*/
|
||||
@Test
|
||||
void onlyPrimaryAndArchitectMaySeeTheMembersArray() {
|
||||
assertTrue(FleetMcp.membersVisibleTo(PRIMARY), "the primary must see the members array");
|
||||
assertTrue(FleetMcp.membersVisibleTo(ARCH_DESIGN), "an architect must see the members array");
|
||||
assertFalse(FleetMcp.membersVisibleTo(WORKER_A),
|
||||
"a worker holds READ but can never SEND, so it must not see the members array");
|
||||
assertFalse(FleetMcp.membersVisibleTo(COLLABORATOR),
|
||||
"a collaborator may SEND to a lead, never to a spawned member, so it must not see the members array");
|
||||
assertFalse(FleetMcp.membersVisibleTo(ANON), "authenticated as nothing must not see it either");
|
||||
}
|
||||
|
||||
/**
|
||||
* Same reasoning as {@link #theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}: the
|
||||
* predicate above can be perfectly correct while the one production call site never asks it.
|
||||
* This reads {@code FleetMcp.java}'s own source and asserts the {@code fleet_list} handler's
|
||||
* {@code listFleet(...)} call asks {@code leadsVisibleTo(principal(exchange))} for the leads
|
||||
* visibility flag, rather than a literal boolean.
|
||||
*/
|
||||
@Test
|
||||
void theFleetListHandlerActuallyConsultsLeadsVisibleTo() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
int start = source.indexOf("listHandler =");
|
||||
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
|
||||
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
|
||||
int end = source.indexOf("stopHandler =", start);
|
||||
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
|
||||
String handlerBlock = source.substring(start, end);
|
||||
|
||||
// CONTROL: the block we scraped really does contain a call to listFleet(...) -- if this
|
||||
// fails, the anchors above moved and the assertion below would otherwise pass on nothing.
|
||||
assertTrue(handlerBlock.contains("listFleet("),
|
||||
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
|
||||
+ "the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("leadsVisibleTo(principal(exchange))"),
|
||||
"the fleet_list handler must ask leadsVisibleTo(principal(exchange)) who is calling, "
|
||||
+ "not pass a literal boolean -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
/** As {@link #theFleetListHandlerActuallyConsultsLeadsVisibleTo}, for {@code membersVisibleTo}. */
|
||||
@Test
|
||||
void theFleetListHandlerActuallyConsultsMembersVisibleTo() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
int start = source.indexOf("listHandler =");
|
||||
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
|
||||
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
|
||||
int end = source.indexOf("stopHandler =", start);
|
||||
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
|
||||
String handlerBlock = source.substring(start, end);
|
||||
|
||||
assertTrue(handlerBlock.contains("listFleet("),
|
||||
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
|
||||
+ "the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("membersVisibleTo(principal(exchange))"),
|
||||
"the fleet_list handler must ask membersVisibleTo(principal(exchange)) who is calling, "
|
||||
+ "not pass a literal boolean -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code fleet_poll{ticket}} must thread the calling connection's own terminal into
|
||||
* {@link MessageService#poll(String, String)}, so a worker cannot read a ticket a different
|
||||
* session created.
|
||||
*/
|
||||
@Test
|
||||
void theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
int start = source.indexOf("pollHandler =");
|
||||
assertTrue(start >= 0, "could not find the fleet_poll handler (pollHandler) in " + MCP_SOURCE
|
||||
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
|
||||
int end = source.indexOf("ackHandler =", start);
|
||||
assertTrue(end > start, "could not find the handler declared after pollHandler to bound the scrape");
|
||||
String handlerBlock = source.substring(start, end);
|
||||
|
||||
// CONTROL: the block we scraped really does call poll(...) -- if this fails, the anchors
|
||||
// above moved and the assertions below would otherwise pass on nothing.
|
||||
assertTrue(handlerBlock.contains("poll(messages,"),
|
||||
"control failed: the scraped pollHandler block contains no poll(messages, ...) call "
|
||||
+ "at all -- the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
|
||||
"the fleet_poll handler must thread callerTerminal(exchange) into poll(...), not omit "
|
||||
+ "it or pass a literal null -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@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) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -122,6 +122,6 @@ class FleetMcpLeadContextGaugeWiringTest {
|
||||
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), contextGauge, leadConfigDirs,
|
||||
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, Map.of(), false,
|
||||
FleetMcp.CoordinationSource.none(), false);
|
||||
FleetMcp.CoordinationSource.none(), false, true, true);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -77,7 +77,7 @@ class FleetMcpTest {
|
||||
|
||||
private void assertSendRoundTrips(String target, Set<String> profiles) throws Exception {
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles));
|
||||
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
@@ -103,7 +103,7 @@ class FleetMcpTest {
|
||||
void sendThenReplyRoundTrips() throws Exception {
|
||||
// fleet_send blocks; fleet_reply resolves it with the worker's structured answer.
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of()));
|
||||
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null));
|
||||
|
||||
// Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now
|
||||
// queues in the inbox if no waiter is open, which would break the round-trip).
|
||||
@@ -181,7 +181,7 @@ class FleetMcpTest {
|
||||
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L));
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
|
||||
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
|
||||
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
@@ -279,7 +279,7 @@ class FleetMcpTest {
|
||||
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L));
|
||||
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null));
|
||||
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
@@ -288,6 +288,89 @@ class FleetMcpTest {
|
||||
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
|
||||
}
|
||||
|
||||
// --- fleetd #715: fleet_send{turnId} is gated on the caller that owns the turn -------------
|
||||
|
||||
/**
|
||||
* A blocking {@code fleet_send} from one caller opens the turn; a {@code fleet_send{turnId}}
|
||||
* from a different caller is refused as an error, and the real owner's answer still succeeds.
|
||||
*/
|
||||
@Test
|
||||
void aDifferentCallersMcpAnswerIsRefusedForABlockingSendButTheRealOwnerSucceeds() throws Exception {
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner"));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T));
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.ask(messages, T, "which config?", 5000L));
|
||||
McpSchema.CallToolResult question = send.get(5, TimeUnit.SECONDS);
|
||||
assertTrue(textOf(question).contains("[question]"), textOf(question));
|
||||
String questionText = textOf(question);
|
||||
String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length());
|
||||
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
|
||||
|
||||
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
|
||||
assertTrue(hijacked.isError(), "a caller that did not open this turn must get an error, not an answer");
|
||||
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
|
||||
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
FleetMcp.reply(messages, T, Role.WORKER, "done");
|
||||
assertEquals("done", textOf(answer.get(5, TimeUnit.SECONDS)));
|
||||
}
|
||||
|
||||
/**
|
||||
* Same hijack and control through the fire-and-poll ({@code sendAsync}) path: the owner comes
|
||||
* from the ticket's recorded creator terminal, not from a caller threaded through a live call.
|
||||
*/
|
||||
@Test
|
||||
void aDifferentCallersMcpAnswerIsRefusedForAnAsyncSendButTheRealOwnerSucceeds() throws Exception {
|
||||
McpSchema.CallToolResult accepted =
|
||||
FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), "term_owner");
|
||||
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
|
||||
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T));
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.ask(messages, T, "which config?", 5000L));
|
||||
MessageService.TaskView asking = messages.poll(ticket);
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
asking = messages.poll(ticket);
|
||||
}
|
||||
assertEquals(MessageService.Phase.ASKING, asking.phase());
|
||||
String turnId = asking.turnId();
|
||||
|
||||
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
|
||||
assertTrue(hijacked.isError(), "a caller that did not create this delegation must get an error");
|
||||
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
|
||||
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
|
||||
"a refused answer must not advance the async ticket's phase");
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
|
||||
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
FleetMcp.reply(messages, T, Role.WORKER, "done");
|
||||
assertEquals("done", textOf(answer.get(5, TimeUnit.SECONDS)));
|
||||
}
|
||||
|
||||
@Test
|
||||
void pollUnknownTicketIsAnError() {
|
||||
McpSchema.CallToolResult res = FleetMcp.poll(messages, "task-999", null);
|
||||
@@ -297,7 +380,7 @@ class FleetMcpTest {
|
||||
|
||||
@Test
|
||||
void sendTimesOutWithAWorkingNote() {
|
||||
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of());
|
||||
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null);
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
|
||||
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
|
||||
}
|
||||
@@ -313,7 +396,7 @@ class FleetMcpTest {
|
||||
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
|
||||
herdr.agentSendFailsWith("send_failed");
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of()));
|
||||
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null));
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
@@ -333,13 +416,13 @@ class FleetMcpTest {
|
||||
|
||||
@Test
|
||||
void sendRejectsMissingArgs() {
|
||||
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError());
|
||||
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of()).isError());
|
||||
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null).isError());
|
||||
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null).isError());
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
|
||||
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"));
|
||||
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null);
|
||||
McpSchema.CallToolResult async = FleetMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
|
||||
|
||||
assertTrue(blocking.isError());
|
||||
@@ -358,7 +441,7 @@ class FleetMcpTest {
|
||||
assertSendRoundTrips("term_live_member", profiles);
|
||||
|
||||
// A herdr-owned pane outside the bridge roster cannot be classified at accept time.
|
||||
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles);
|
||||
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null);
|
||||
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
|
||||
}
|
||||
|
||||
@@ -450,7 +533,7 @@ class FleetMcpTest {
|
||||
void askThenAnswerRoundTrips() throws Exception {
|
||||
// The primary delegates and blocks; wait until its waiter is open before the worker asks.
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of()));
|
||||
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
@@ -472,7 +555,7 @@ class FleetMcpTest {
|
||||
|
||||
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L));
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
|
||||
|
||||
// The worker's ask returns the answer — it resumes the same turn.
|
||||
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
|
||||
@@ -498,7 +581,7 @@ class FleetMcpTest {
|
||||
|
||||
@Test
|
||||
void answerToAStaleTurnIsAnError() {
|
||||
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L);
|
||||
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null);
|
||||
assertTrue(res.isError());
|
||||
assertTrue(textOf(res).contains("no longer open"), textOf(res));
|
||||
}
|
||||
@@ -692,7 +775,7 @@ class FleetMcpTest {
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
|
||||
|
||||
String out = textOf(res);
|
||||
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
|
||||
@@ -721,7 +804,7 @@ class FleetMcpTest {
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out);
|
||||
@@ -780,19 +863,82 @@ class FleetMcpTest {
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
boolean callerIsPrimary = Principal.worker("term_a", 1).isPrimary();
|
||||
Principal worker = Principal.worker("term_a", 1);
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), worker.isPrimary(),
|
||||
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
|
||||
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("\"coordinator\""), "a worker must never see the coordinator key at all: " + out);
|
||||
assertFalse(out.contains("mac-opus"), "no fragment of the coordinator row may leak either: " + out);
|
||||
assertTrue(out.contains("\"leads\""), "the rest of the result must still be present: " + out);
|
||||
assertTrue(out.contains("\"members\""), out);
|
||||
assertFalse(out.contains("\"leads\""), "a worker must never see the leads key at all: " + out);
|
||||
assertFalse(out.contains("\"members\""), "a worker must never see the members key at all: " + out);
|
||||
assertTrue(out.contains("\"healthCoverage\""), "the rest of the result must still be present: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* A worker's result must carry neither the {@code leads} nor the {@code members} key, and no
|
||||
* fragment of either row leaks even though both are fully populated for this call -- a key
|
||||
* check alone would pass on an implementation that still built the rows and only renamed or
|
||||
* nested them.
|
||||
*/
|
||||
@Test
|
||||
void listLeaksNoLeadOrMemberRowFragmentToAWorker() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_owner",
|
||||
new WorktreeRequest("cb-304", null));
|
||||
Principal worker = Principal.worker("term_a", 1);
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of("term_lead_x", "mac-opus"), "", FleetMcp.CoordinationSource.none(), worker.isPrimary(),
|
||||
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
|
||||
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("\"leads\""), out);
|
||||
assertFalse(out.contains("\"members\""), out);
|
||||
assertFalse(out.contains(s.terminalId()), "no member row fragment may leak to a worker: " + out);
|
||||
assertFalse(out.contains("term_owner"), "no member owner fragment may leak to a worker: " + out);
|
||||
assertFalse(out.contains("mac-opus"), "no lead row fragment may leak to a worker: " + out);
|
||||
assertTrue(out.contains("\"healthCoverage\""), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* A collaborator may {@code SEND} to a lead, so it must see the {@code leads} array -- the
|
||||
* only place {@code fleet_whoami} does not already give it a lead's address. It may never
|
||||
* {@code SEND} to a spawned member, so the {@code members} key must stay absent for it, with
|
||||
* no fragment of a populated member row leaking either.
|
||||
*/
|
||||
@Test
|
||||
void listShowsLeadsButNotMembersToACollaborator() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_owner",
|
||||
new WorktreeRequest("cb-304", null));
|
||||
Principal collaborator = Principal.collaborator("ops", "term_collab", 600);
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of("term_lead_x", "mac-opus"), "", FleetMcp.CoordinationSource.none(), collaborator.isPrimary(),
|
||||
FleetMcp.leadsVisibleTo(collaborator), FleetMcp.membersVisibleTo(collaborator));
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"leads\""), "a collaborator must see the leads array: " + out);
|
||||
assertTrue(out.contains("mac-opus"), "a collaborator must see the lead's name/address: " + out);
|
||||
assertFalse(out.contains("\"members\""), "a collaborator must never see the members key: " + out);
|
||||
assertFalse(out.contains(s.terminalId()), "no member row fragment may leak to a collaborator: " + out);
|
||||
assertFalse(out.contains("term_owner"), "no member owner fragment may leak to a collaborator: " + out);
|
||||
assertTrue(out.contains("\"healthCoverage\""), out);
|
||||
}
|
||||
|
||||
@@ -808,16 +954,19 @@ class FleetMcpTest {
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
boolean callerIsPrimary = Principal.architect("lead-designer", "term_design", 400).isPrimary();
|
||||
Principal architect = Principal.architect("lead-designer", "term_design", 400);
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), architect.isPrimary(),
|
||||
FleetMcp.leadsVisibleTo(architect), FleetMcp.membersVisibleTo(architect));
|
||||
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("\"coordinator\""), "an architect must never see the coordinator key either: " + out);
|
||||
assertTrue(out.contains("\"leads\""), "an architect must still see the leads array: " + out);
|
||||
assertTrue(out.contains("\"members\""), "an architect must still see the members array: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -842,7 +991,7 @@ class FleetMcpTest {
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true));
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true));
|
||||
|
||||
assertTrue(gatedAsPrimary.contains("\"coordinator\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"selfId\":\"mac-opus\""), gatedAsPrimary);
|
||||
@@ -851,6 +1000,8 @@ class FleetMcpTest {
|
||||
assertTrue(gatedAsPrimary.contains("\"heldDurable\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"held\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"peers\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"leads\""), "the primary must still see the leads array: " + gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"members\""), "the primary must still see the members array: " + gatedAsPrimary);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -869,7 +1020,7 @@ class FleetMcpTest {
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"msgId\":\"m1\""), out);
|
||||
@@ -893,7 +1044,7 @@ class FleetMcpTest {
|
||||
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
|
||||
Map.of(), "", collaborators, collaboratorsVisible,
|
||||
FleetMcp.CoordinationSource.none(), false);
|
||||
FleetMcp.CoordinationSource.none(), false, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -978,7 +1129,7 @@ class FleetMcpTest {
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"pending\":0"), out);
|
||||
@@ -1007,7 +1158,7 @@ class FleetMcpTest {
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"heldDurable\":false"),
|
||||
@@ -1076,7 +1227,7 @@ class FleetMcpTest {
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true);
|
||||
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true, true, true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out);
|
||||
@@ -1792,14 +1943,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
|
||||
@@ -1822,7 +1973,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);
|
||||
@@ -1833,7 +1984,64 @@ class FleetMcpTest {
|
||||
// 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));
|
||||
() -> messages.answer(turnId, "config.yaml", 5000, null));
|
||||
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);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@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, "term_creator"));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -78,7 +78,7 @@ class MessageServiceTest {
|
||||
}
|
||||
|
||||
private CompletableFuture<MessageService.Reply> sendAsync(String content, long timeoutMillis) {
|
||||
return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis));
|
||||
return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis, null));
|
||||
}
|
||||
|
||||
private void awaitWaiting() throws InterruptedException {
|
||||
@@ -188,7 +188,7 @@ class MessageServiceTest {
|
||||
localRendezvous, new InMemoryReplyInbox());
|
||||
String brief = "Implement the requested change. ".repeat(20);
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000));
|
||||
CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000, (String) null));
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!localRendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
@@ -341,7 +341,7 @@ class MessageServiceTest {
|
||||
|
||||
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null));
|
||||
|
||||
// The worker's ask returns the answer — it resumes the same turn.
|
||||
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
|
||||
@@ -377,7 +377,7 @@ class MessageServiceTest {
|
||||
|
||||
// The primary answers that one turnId; both asks unblock with the same answer.
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null));
|
||||
|
||||
MessageService.AskResult a1 = ask1.get(5, TimeUnit.SECONDS);
|
||||
MessageService.AskResult a2 = ask2.get(5, TimeUnit.SECONDS);
|
||||
@@ -472,7 +472,7 @@ class MessageServiceTest {
|
||||
"the duplicate's own timeout elapses first");
|
||||
|
||||
CompletableFuture<MessageService.Reply> answered =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500));
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500, null));
|
||||
|
||||
MessageService.AskResult a = fresh.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome(),
|
||||
@@ -507,7 +507,7 @@ class MessageServiceTest {
|
||||
|
||||
CompletableFuture<MessageService.Reply> lateAnswer = new CompletableFuture<>();
|
||||
messages.setAskTimeoutRaceHookForTest(() ->
|
||||
lateAnswer.complete(messages.answer(turnId, "too late", 500)));
|
||||
lateAnswer.complete(messages.answer(turnId, "too late", 500, null)));
|
||||
try {
|
||||
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome());
|
||||
@@ -539,7 +539,7 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void answeringAnUnknownTurnIsStale() {
|
||||
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
|
||||
MessageService.Reply r = messages.answer(T + "#999", "too late", 500, null);
|
||||
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
|
||||
"an answer to a turn that never existed (or already lapsed) is stale, not a hang");
|
||||
}
|
||||
@@ -572,7 +572,7 @@ class MessageServiceTest {
|
||||
|
||||
messages.setAnswerAskLapseRaceHookForTest(() -> rendezvous.answerAsk(turnId, "raced in first"));
|
||||
try {
|
||||
MessageService.Reply r = messages.answer(turnId, "too late", 500);
|
||||
MessageService.Reply r = messages.answer(turnId, "too late", 500, null);
|
||||
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
|
||||
"an ask already answered by the race must be seen as lapsed, not double-delivered");
|
||||
} finally {
|
||||
@@ -596,7 +596,7 @@ class MessageServiceTest {
|
||||
void sendTimesOutBeforeDeliveryIsQueuedNotWorking() {
|
||||
// Nothing ever delivers the message and nothing resolves the send, so the reply future
|
||||
// times out with delivery still incomplete — the message is still queued for the worker.
|
||||
MessageService.Reply r = messages.send(T, "never delivered", 50);
|
||||
MessageService.Reply r = messages.send(T, "never delivered", 50, null);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome(),
|
||||
"an undelivered send that times out is still queued, not working");
|
||||
assertNull(r.text());
|
||||
@@ -605,7 +605,7 @@ class MessageServiceTest {
|
||||
@Test
|
||||
void sendTimesOutAfterDeliveryIsStillWorking() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 300));
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 300, null));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver — the delivered future now completes
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker starts but never replies
|
||||
@@ -630,7 +630,7 @@ class MessageServiceTest {
|
||||
void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
|
||||
messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
|
||||
try {
|
||||
MessageService.Reply reply = messages.send(T, "race delivery", 50);
|
||||
MessageService.Reply reply = messages.send(T, "race delivery", 50, null);
|
||||
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
|
||||
"cancel reporting DELIVERED means the worker received the timed-out message");
|
||||
@@ -652,7 +652,7 @@ class MessageServiceTest {
|
||||
void sendTimesOutWithAttemptedDeliveryReportsUnconfirmedNotQueued() throws Exception {
|
||||
herdr.agentSendFailsWith("send_failed");
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150));
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150, null));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt → ATTEMPTED
|
||||
|
||||
@@ -679,7 +679,7 @@ class MessageServiceTest {
|
||||
|
||||
// The primary answers, unblocking the worker; but the worker never sends the follow-up
|
||||
// fleet_reply, so the answering send rides out its short window as still-working.
|
||||
MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200);
|
||||
MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200, null);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answer.outcome(),
|
||||
"an answered worker that never replies times out as still working");
|
||||
|
||||
@@ -715,7 +715,7 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null));
|
||||
ask.get(5, TimeUnit.SECONDS); // worker resumed with the answer
|
||||
|
||||
awaitWaiting(); // the answering call has (re)opened its own forward waiter
|
||||
@@ -723,7 +723,7 @@ class MessageServiceTest {
|
||||
MessageService.Reply done = answer.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.REPLIED, done.outcome());
|
||||
|
||||
MessageService.Reply probe = messages.send(T, "probe after normal reply", 300);
|
||||
MessageService.Reply probe = messages.send(T, "probe after normal reply", 300, null);
|
||||
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
|
||||
"the session lock must be released after a normal REPLIED answer(), or this bounded "
|
||||
+ "follow-up send would come back BUSY instead of timing out on its own work");
|
||||
@@ -756,13 +756,13 @@ class MessageServiceTest {
|
||||
// The primary answers, unblocking the worker; the worker never sends its follow-up
|
||||
// fleet_reply, so the answering call rides out its short window as still-working.
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200));
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200, null));
|
||||
MessageService.Reply answered = answer.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answered.outcome(),
|
||||
"an answered worker that never replies times out as still working");
|
||||
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
|
||||
|
||||
MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300);
|
||||
MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300, null);
|
||||
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
|
||||
"the session lock must be released after a TIMED_OUT_WORKING answer(), or this bounded "
|
||||
+ "follow-up send would come back BUSY instead of timing out on its own work");
|
||||
@@ -791,7 +791,7 @@ class MessageServiceTest {
|
||||
new java.util.concurrent.atomic.AtomicReference<>();
|
||||
Thread answerer = new Thread(() -> {
|
||||
try {
|
||||
messages.answer(turnId, "config.yaml", 5000);
|
||||
messages.answer(turnId, "config.yaml", 5000, null);
|
||||
caught.set(new AssertionError("expected answer() to throw"));
|
||||
} catch (Throwable t) {
|
||||
caught.set(t);
|
||||
@@ -813,7 +813,7 @@ class MessageServiceTest {
|
||||
|
||||
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
|
||||
|
||||
MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300);
|
||||
MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300, null);
|
||||
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
|
||||
"the session lock must be released when answer()'s reply future fails exceptionally, "
|
||||
+ "or this bounded follow-up send would come back BUSY instead of timing out on its own work");
|
||||
@@ -840,7 +840,7 @@ class MessageServiceTest {
|
||||
new java.util.concurrent.atomic.AtomicReference<>();
|
||||
Thread answerer = new Thread(() -> {
|
||||
try {
|
||||
messages.answer(turnId, "config.yaml", 5000);
|
||||
messages.answer(turnId, "config.yaml", 5000, null);
|
||||
caught.set(new AssertionError("expected answer() to throw"));
|
||||
} catch (Throwable t) {
|
||||
caught.set(t);
|
||||
@@ -859,17 +859,88 @@ class MessageServiceTest {
|
||||
|
||||
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
|
||||
|
||||
MessageService.Reply probe = messages.send(T, "probe after interruption", 300);
|
||||
MessageService.Reply probe = messages.send(T, "probe after interruption", 300, null);
|
||||
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
|
||||
"the session lock must be released when answer()'s wait is interrupted, or this bounded "
|
||||
+ "follow-up send would come back BUSY instead of timing out on its own work");
|
||||
}
|
||||
|
||||
/**
|
||||
* Drive an async send on {@code T} to a resolved reply, polling as {@code owner} until the
|
||||
* ticket reports {@link MessageService.Phase#DONE} (or the 2s deadline runs out). {@code owner}
|
||||
* must be a terminal this ticket's creator check actually accepts, or this loops to the
|
||||
* deadline and returns a non-{@code DONE} view.
|
||||
*/
|
||||
private MessageService.TaskView driveAsyncTicketToDone(String ticket, String owner) throws InterruptedException {
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (view == null || view.phase() != MessageService.Phase.DONE) {
|
||||
if (System.currentTimeMillis() >= deadline) break;
|
||||
view = messages.poll(ticket, owner);
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
return view;
|
||||
}
|
||||
|
||||
@Test
|
||||
void pollReturnsNullForAnUnknownTicket() {
|
||||
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
|
||||
}
|
||||
|
||||
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
|
||||
|
||||
@Test
|
||||
void pollByAnotherTerminalIsRefused() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_a");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "secret async result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, "term_a");
|
||||
assertNotNull(owner, "the creator must still be able to read its own ticket");
|
||||
assertEquals(MessageService.Phase.DONE, owner.phase());
|
||||
|
||||
MessageService.TaskView refused = messages.poll(ticket, "term_b");
|
||||
assertNotNull(refused, "a different terminal gets a refusal, not silence");
|
||||
assertNotEquals(MessageService.Phase.DONE, refused.phase(),
|
||||
"a different terminal must never see the ticket as DONE");
|
||||
assertNull(refused.reply(), "a refusal must never carry the reply text");
|
||||
assertFalse(String.valueOf(refused).contains("secret async result"),
|
||||
"the reply text must not appear anywhere in the refused view");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
|
||||
|
||||
// callerTerminal == null is the unnamed primary (resolved by token or loopback trust, with
|
||||
// no herdr pane) — it must read a ticket a terminal-bearing lead created.
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
|
||||
assertNotNull(view, "the unnamed primary must be able to read any ticket");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("primary-visible result", view.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void creatorReadsItsOwnTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "own result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, "term_creator");
|
||||
assertNotNull(view, "the ticket's own creator must be able to read it");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("own result", view.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void pollReportsACompletedTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
@@ -896,11 +967,11 @@ class MessageServiceTest {
|
||||
@Test
|
||||
void concurrentSendToSameSessionWhileFirstHoldsItIsBusy() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> first =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000));
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000, null));
|
||||
awaitWaiting(); // the first send now holds the session lock, blocked on its reply
|
||||
|
||||
// A second send to the SAME session cannot take the lock within its short window.
|
||||
MessageService.Reply busy = messages.send(T, "second", 100);
|
||||
MessageService.Reply busy = messages.send(T, "second", 100, null);
|
||||
assertEquals(MessageService.Outcome.BUSY, busy.outcome(),
|
||||
"a second send while another holds the session is busy, not a hang");
|
||||
assertNull(busy.text());
|
||||
@@ -930,13 +1001,13 @@ class MessageServiceTest {
|
||||
PrimaryRegistry reg = new PrimaryRegistry(null);
|
||||
// L accepts a delegation to W: the send wins the lock and queues delivery → L is recorded.
|
||||
CompletableFuture<MessageService.Reply> first = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L)));
|
||||
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L), null));
|
||||
awaitWaiting();
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
"an accepted send owns the delegation");
|
||||
|
||||
// A attempts W while L holds it → BUSY (lock never taken) → its hook never fires.
|
||||
MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A));
|
||||
MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A), null);
|
||||
assertEquals(MessageService.Outcome.BUSY, busy.outcome());
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
"a BUSY send must not steal the delegator ownership it never earned");
|
||||
@@ -957,7 +1028,7 @@ class MessageServiceTest {
|
||||
void anAcceptedSendAfterThePriorOwnerFinishesBecomesTheNewOwner() throws Exception {
|
||||
PrimaryRegistry reg = new PrimaryRegistry(null);
|
||||
CompletableFuture<MessageService.Reply> first = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L)));
|
||||
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L), null));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
@@ -968,7 +1039,7 @@ class MessageServiceTest {
|
||||
|
||||
// L finished; A's later accepted send takes the delegation over.
|
||||
CompletableFuture<MessageService.Reply> second = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A)));
|
||||
() -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A), null));
|
||||
awaitWaiting();
|
||||
assertEquals(LEAD_A, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
"an accepted send after the owner finished becomes the new delegator");
|
||||
@@ -989,7 +1060,7 @@ class MessageServiceTest {
|
||||
void answeringAnAskDoesNotRewriteDelegatorOwnership() throws Exception {
|
||||
PrimaryRegistry reg = new PrimaryRegistry(null);
|
||||
CompletableFuture<MessageService.Reply> send = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L)));
|
||||
() -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L), null));
|
||||
awaitWaiting();
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "L owns the delegation");
|
||||
|
||||
@@ -1004,7 +1075,7 @@ class MessageServiceTest {
|
||||
|
||||
// L answers the ask on the same turn; the answer path must not touch ownership.
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting(); // the answering send reopened its forward waiter
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
@@ -1024,7 +1095,7 @@ class MessageServiceTest {
|
||||
void aThrowingAcceptedHookLeavesNoStaleWaiterOrQueuedOrphan() {
|
||||
assertThrows(IllegalStateException.class,
|
||||
() -> messages.send(T, "doomed", 500,
|
||||
() -> { throw new IllegalStateException("ownership hook failed"); }),
|
||||
() -> { throw new IllegalStateException("ownership hook failed"); }, null),
|
||||
"a throwing ownership hook fails the send loudly");
|
||||
|
||||
assertFalse(rendezvous.isWaiting(T), "the failed send must not leave a stale rendezvous waiter");
|
||||
@@ -1152,7 +1223,7 @@ class MessageServiceTest {
|
||||
@Test
|
||||
void abandonFailsASendThatIsStillWaitingOnAReleasedSession() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000, null));
|
||||
awaitWaiting();
|
||||
|
||||
assertTrue(messages.abandon(T, "session released"), "a live waiter is abandoned");
|
||||
@@ -1172,7 +1243,7 @@ class MessageServiceTest {
|
||||
@Test
|
||||
void abandonDoesNotOverwriteAnAlreadyResolvedSend() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000, null));
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "the real answer"));
|
||||
|
||||
@@ -1303,7 +1374,7 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, null));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
@@ -1343,7 +1414,7 @@ class MessageServiceTest {
|
||||
// same turnId must see it as lapsed rather than resolving a question nobody is waiting on.
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT, ask.get(5, TimeUnit.SECONDS).outcome());
|
||||
assertEquals(MessageService.Outcome.STALE_TURN,
|
||||
messages.answer(asking.turnId(), "config.yaml", 200).outcome());
|
||||
messages.answer(asking.turnId(), "config.yaml", 200, null).outcome());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -1380,7 +1451,7 @@ class MessageServiceTest {
|
||||
|
||||
// The primary answers, but its own bounded wait for the worker's resumed turn is short and
|
||||
// expires before the worker (still genuinely working) gets back to it.
|
||||
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
|
||||
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150, null);
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
|
||||
"the primary's own bounded wait gives up before the worker finishes resuming");
|
||||
@@ -1409,7 +1480,7 @@ class MessageServiceTest {
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
|
||||
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150, null);
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
|
||||
|
||||
@@ -1463,7 +1534,7 @@ class MessageServiceTest {
|
||||
messages.setFinishAsyncTaskRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000, null));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
|
||||
|
||||
@@ -1569,7 +1640,7 @@ class MessageServiceTest {
|
||||
String turnId = asking.turnId();
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000, null));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer(),
|
||||
"the worker's own ask() call must have already unblocked with the primary's answer "
|
||||
+ "before we force the race below");
|
||||
@@ -1616,7 +1687,7 @@ class MessageServiceTest {
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
String turnId = asking.turnId();
|
||||
|
||||
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150);
|
||||
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150, null);
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
|
||||
"the primary's own bounded wait must give up first, leaving turnId stamped with no "
|
||||
@@ -1682,7 +1753,7 @@ class MessageServiceTest {
|
||||
|
||||
MessageService.TaskView asking = messages.poll(first);
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, null));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
@@ -1716,7 +1787,7 @@ class MessageServiceTest {
|
||||
|
||||
// The primary answers it — answer() resumes the turn and blocks for what comes next.
|
||||
CompletableFuture<MessageService.Reply> answer1 = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking1.turnId(), "a1", 5000));
|
||||
() -> messages.answer(asking1.turnId(), "a1", 5000, null));
|
||||
assertEquals("a1", ask1.get(5, TimeUnit.SECONDS).answer());
|
||||
|
||||
// Still in the SAME resumed turn — before replying — the worker asks again.
|
||||
@@ -1737,7 +1808,7 @@ class MessageServiceTest {
|
||||
|
||||
// The primary answers the second question; the worker finally sends its real fleet_reply.
|
||||
CompletableFuture<MessageService.Reply> answer2 = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(turnId2, "a2", 5000));
|
||||
() -> messages.answer(turnId2, "a2", 5000, null));
|
||||
assertEquals("a2", ask2.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
@@ -1804,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
|
||||
@@ -1830,21 +1901,151 @@ 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());
|
||||
assertEquals(asking.turnId(), pending.turnId());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, null));
|
||||
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, "term_creator"));
|
||||
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, null));
|
||||
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());
|
||||
}
|
||||
|
||||
// --- fleetd #715: answer() is gated on the caller that owns the turn -----------------------
|
||||
|
||||
/**
|
||||
* A blocking {@code fleet_send} from one caller opens the turn; a different caller's answer is
|
||||
* refused with no side effect on the rendezvous or the worker's blocked {@code fleet_ask} —
|
||||
* only the real owner can answer it.
|
||||
*/
|
||||
@Test
|
||||
void aDifferentCallersAnswerIsRefusedForABlockingSendDelegationButTheRealOwnerSucceeds() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "do X", 5000, "term_owner"));
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
||||
assertFalse(rendezvous.isWaiting(T), "the forward waiter closes once the question surfaces");
|
||||
|
||||
// The hijack: a different caller answers the SAME turnId.
|
||||
MessageService.Reply hijacked = messages.answer(q.turnId(), "evil.yaml", 500, "term_attacker");
|
||||
assertEquals(MessageService.Outcome.NOT_TURN_OWNER, hijacked.outcome(),
|
||||
"a caller that did not open this turn must be refused, not served");
|
||||
assertFalse(rendezvous.isWaiting(T), "a refused answer must not open a forward waiter");
|
||||
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
|
||||
assertEquals(T, rendezvous.askSession(q.turnId()), "a refused answer must leave the ask turn open");
|
||||
|
||||
// The control: the real owner answers the same turnId and the turn resumes normally.
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(q.turnId(), "config.yaml", 5000, "term_owner"));
|
||||
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());
|
||||
}
|
||||
|
||||
/**
|
||||
* Same hijack and control as the blocking case, but the delegation is opened through the
|
||||
* fire-and-poll path ({@code sendAsync}) — the owner comes from the ticket's recorded creator
|
||||
* terminal, not a caller argument threaded through a live blocking call.
|
||||
*/
|
||||
@Test
|
||||
void aDifferentCallersAnswerIsRefusedForAnAsyncSendDelegationButTheRealOwnerSucceeds() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, "term_owner");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.Reply hijacked = messages.answer(asking.turnId(), "evil.yaml", 500, "term_attacker");
|
||||
assertEquals(MessageService.Outcome.NOT_TURN_OWNER, hijacked.outcome(),
|
||||
"a caller that did not create this delegation must be refused, not served");
|
||||
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
|
||||
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
|
||||
"a refused answer must not advance the async ticket's phase");
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_owner"));
|
||||
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());
|
||||
assertEquals("done", awaitTicketPhase(ticket, MessageService.Phase.DONE).reply());
|
||||
}
|
||||
|
||||
// --- CB-588: async ticket terminal nudges ---------------------------------------------------
|
||||
//
|
||||
// MessageService.reply's rendezvous fast path is exactly what an async ticket always takes
|
||||
@@ -2101,7 +2302,7 @@ class MessageServiceTest {
|
||||
|
||||
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000, null));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
@@ -2127,7 +2328,7 @@ class MessageServiceTest {
|
||||
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000, null));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
@@ -2520,7 +2721,7 @@ class MessageServiceTest {
|
||||
void hasQueuedDeliveryIsTrueAfterAnUndeliveredSendTimesOut() {
|
||||
// Nothing ever delivers the message (never goes IDLE/BLOCKED), so the send times out with
|
||||
// TIMED_OUT_QUEUED — same setup as sendTimesOutBeforeDeliveryIsQueuedNotWorking above.
|
||||
MessageService.Reply r = messages.send(T, "never delivered", 50);
|
||||
MessageService.Reply r = messages.send(T, "never delivered", 50, null);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
|
||||
|
||||
assertTrue(messages.hasQueuedDelivery(T),
|
||||
@@ -2532,7 +2733,7 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50, null).outcome());
|
||||
assertTrue(messages.hasQueuedDelivery(T));
|
||||
|
||||
// A fresh send accepts delivery (opens its own waiter) — the stale queued fact is cleared.
|
||||
@@ -2549,7 +2750,7 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryClearsOnAbandon() {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50, null).outcome());
|
||||
assertTrue(messages.hasQueuedDelivery(T));
|
||||
|
||||
messages.abandon(T, "session released");
|
||||
|
||||
@@ -0,0 +1,187 @@
|
||||
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.stream.Stream;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* Pins that no file under {@code src/main/java} calls the fail-open
|
||||
* {@link Rendezvous#open(String)} overload. That overload opens a forward waiter with no
|
||||
* recorded {@link Rendezvous.Owner}, so the turn it opens can never be answered by anyone,
|
||||
* not even the unnamed primary. Every production caller must go through
|
||||
* {@link Rendezvous#open(String, Rendezvous.Owner)} and record an explicit owner.
|
||||
*
|
||||
* <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 rendezvous.open(}, classified by
|
||||
* argument count. The one test method here points that exact scanner at a file known to hold
|
||||
* many real one-argument calls before it ever looks at production, so a scanner that stops
|
||||
* matching fails loudly on the known-positive case instead of leaving a clean production result
|
||||
* looking like evidence it never produced.
|
||||
*/
|
||||
class RendezvousOpenUsageTest {
|
||||
|
||||
private static final Path PRODUCTION_SOURCE = Path.of("src/main/java");
|
||||
private static final Path KNOWN_TEST_CALLER =
|
||||
Path.of("src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java");
|
||||
|
||||
/**
|
||||
* A pattern that cannot find a known one-argument {@code rendezvous.open(} call would also
|
||||
* find none in production -- not because production is clean, but because the pattern does
|
||||
* not match the text it is supposed to catch. That failure mode is exactly what let a
|
||||
* {@code \b}-based regex read "no callers" under {@code git grep -E} when 81 real ones
|
||||
* existed: {@code git grep} does not treat {@code \b} as a word boundary, so the pattern
|
||||
* silently matched nothing anywhere, clean code and real calls alike. This test runs the
|
||||
* known-positive check first, with the same scanning method the production check then
|
||||
* depends on, so that mistake fails loudly here instead of reading as a clean result.
|
||||
*/
|
||||
@Test
|
||||
void noProductionFileCallsTheSingleArgumentOpenOverload() throws IOException {
|
||||
List<String> knownCalls = new ArrayList<>();
|
||||
int knownFilesScanned = scanForOneArgOpenCalls(KNOWN_TEST_CALLER, knownCalls);
|
||||
|
||||
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "found
|
||||
// nothing" having looked at nothing.
|
||||
assertTrue(knownFilesScanned > 0, "control failed: the scan under " + KNOWN_TEST_CALLER
|
||||
+ " visited zero .java files -- the path is wrong, so neither result below proves "
|
||||
+ "anything");
|
||||
|
||||
// CONTROL: the scanner actually finds real one-argument rendezvous.open( calls when
|
||||
// pointed at a file known to hold many. If this is not satisfied, the matching logic
|
||||
// itself is broken, and the production result below is the scanner failing silently,
|
||||
// not production code actually being clean.
|
||||
assertTrue(knownCalls.size() >= 60, "control failed: the scanner found only "
|
||||
+ knownCalls.size() + " one-argument rendezvous.open( call(s) in " + KNOWN_TEST_CALLER
|
||||
+ ", which is known to hold many -- the matching logic itself is broken: " + knownCalls);
|
||||
|
||||
List<String> violations = new ArrayList<>();
|
||||
int filesScanned = scanForOneArgOpenCalls(PRODUCTION_SOURCE, violations);
|
||||
|
||||
// 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 Rendezvous.open(String) "
|
||||
+ "overload, which opens a forward waiter with no recorded owner -- record an "
|
||||
+ "explicit Rendezvous.Owner through open(String, Owner) instead: " + violations);
|
||||
}
|
||||
|
||||
private static int scanForOneArgOpenCalls(Path root, List<String> sites) throws IOException {
|
||||
List<Path> files = javaFiles(root);
|
||||
for (Path file : files) {
|
||||
scanFileForOpenCalls(file, sites);
|
||||
}
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
||||
private static void scanFileForOpenCalls(Path file, List<String> sites) throws IOException {
|
||||
String source = Files.readString(file);
|
||||
String needle = "rendezvous.open(";
|
||||
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; // Rendezvous has no zero-argument open() -- a javadoc "rendezvous.open()" mention, not a call
|
||||
}
|
||||
if (topLevelCommaCount(args) == 0) {
|
||||
sites.add(file + ":" + lineOf(source, idx) + " -- rendezvous.open(" + args.trim() + ")");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The text between {@code rendezvous.open(} 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;
|
||||
}
|
||||
}
|
||||
@@ -132,4 +132,80 @@ class RendezvousTest {
|
||||
assertNull(rendezvous.askSession(t.turnId()), "a closed ask is forgotten");
|
||||
assertFalse(rendezvous.answerAsk(t.turnId(), "late"), "a closed ask can no longer be answered");
|
||||
}
|
||||
|
||||
// ── fleetd #715: turn ownership ────────────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void openWithNoOwnerRecordsNoOwnerAndOpenWithAnOwnerRecordsIt() {
|
||||
rendezvous.open(W);
|
||||
assertNull(rendezvous.ownerOf(W), "the no-owner overload records no owner at all");
|
||||
rendezvous.close(W, rendezvous.currentWaiter(W));
|
||||
|
||||
rendezvous.open(W, Rendezvous.Owner.of("term_lead"));
|
||||
assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.ownerOf(W));
|
||||
}
|
||||
|
||||
@Test
|
||||
void ownerOfIsNullWhenNoWaiterIsOpen() {
|
||||
assertNull(rendezvous.ownerOf(W), "no waiter open means no owner to report");
|
||||
}
|
||||
|
||||
@Test
|
||||
void openAskStampsTheForwardWaitersOwnerOntoTheFreshTurnOnly() {
|
||||
rendezvous.open(W, Rendezvous.Owner.of("term_lead"));
|
||||
|
||||
Rendezvous.AskTicket fresh = rendezvous.openAsk(W);
|
||||
assertTrue(fresh.fresh());
|
||||
assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.askOwner(fresh.turnId()),
|
||||
"a freshly-opened ask copies the forward waiter's current owner");
|
||||
|
||||
Rendezvous.AskTicket coalesced = rendezvous.openAsk(W);
|
||||
assertFalse(coalesced.fresh());
|
||||
assertEquals(fresh.turnId(), coalesced.turnId());
|
||||
assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.askOwner(coalesced.turnId()),
|
||||
"a coalesced duplicate ask rides the fresh owner's turn, unchanged");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSecondAskAfterTheFirstClosesStampsWhateverOwnerIsOpenAtThatLaterMoment() {
|
||||
rendezvous.open(W, Rendezvous.Owner.of("term_lead"));
|
||||
Rendezvous.AskTicket first = rendezvous.openAsk(W);
|
||||
rendezvous.closeAsk(first.turnId());
|
||||
|
||||
// The forward waiter is reopened under a different owner before the second ask — mirrors
|
||||
// answer() reopening with the owner it already checked, which can differ turn to turn.
|
||||
rendezvous.close(W, rendezvous.currentWaiter(W));
|
||||
rendezvous.open(W, Rendezvous.Owner.of("term_other"));
|
||||
|
||||
Rendezvous.AskTicket second = rendezvous.openAsk(W);
|
||||
assertTrue(second.fresh());
|
||||
assertEquals(Rendezvous.Owner.of("term_other"), rendezvous.askOwner(second.turnId()),
|
||||
"a freshly-opened ask always copies whatever owner is open right now, not a stale one");
|
||||
}
|
||||
|
||||
@Test
|
||||
void askOwnerIsNullForAnUnknownOrLapsedTurn() {
|
||||
assertNull(rendezvous.askOwner("no-such#1"));
|
||||
Rendezvous.AskTicket t = rendezvous.openAsk(W);
|
||||
rendezvous.closeAsk(t.turnId());
|
||||
assertNull(rendezvous.askOwner(t.turnId()), "a closed ask no longer reports an owner");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ownerPermitsIsFailClosedOnARecordAndThreeStatesAreDistinct() {
|
||||
assertFalse(Rendezvous.Owner.permits(null, null),
|
||||
"no owner on record refuses even a caller with no terminal");
|
||||
assertFalse(Rendezvous.Owner.permits(null, "term_a"),
|
||||
"no owner on record refuses a terminal-bearing caller too");
|
||||
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, null),
|
||||
"the unnamed primary owner matches a caller with no terminal");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, "term_a"),
|
||||
"the unnamed primary owner does not match a terminal-bearing caller");
|
||||
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_a"),
|
||||
"a named owner matches the same terminal");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_b"),
|
||||
"a named owner refuses a different terminal");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), null),
|
||||
"a named owner refuses the unnamed primary");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,405 @@ 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, "term_a"));
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code POST /sessions/{id}/message} with a {@code turnId} refuses a caller whose terminal
|
||||
* did not open the turn, even when that caller otherwise holds ANSWER rights, and leaves the
|
||||
* turn open for the real owner to resolve. Covers the REST adapter's ANSWER gate for an
|
||||
* async-send ({@code wait:false}) delegation.
|
||||
*/
|
||||
@Test
|
||||
void restAnswerIsRefusedForADifferentCallerOnAnAsyncSendDelegationButTheRealOwnerSucceeds() 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 ownerApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID, Map.of("term_a", "lead-owner"));
|
||||
Javalin attackerApp = startOnSharedService(messages, herdr, 9001L, Map.of("term_shell", "lead-attacker"));
|
||||
try {
|
||||
ObjectMapper mapper = new ObjectMapper();
|
||||
HttpResponse<String> created = send(ownerApp.port(), "POST", "/sessions/term_target/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());
|
||||
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_target"), "the async send 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();
|
||||
|
||||
HttpResponse<String> hijacked = send(attackerApp.port(), "POST", "/sessions/term_target/message",
|
||||
"{\"content\":\"hijack\",\"turnId\":\"" + turnId + "\"}", null);
|
||||
assertEquals(403, hijacked.statusCode(), hijacked.body());
|
||||
assertTrue(hijacked.body().contains("not_turn_owner"),
|
||||
"a different caller's answer must be refused as not_turn_owner: " + hijacked.body());
|
||||
assertEquals("term_target", rendezvous.askSession(turnId),
|
||||
"a refused answer must leave the ask turn open");
|
||||
|
||||
CompletableFuture<HttpResponse<String>> owned = CompletableFuture.supplyAsync(() -> {
|
||||
try {
|
||||
return send(ownerApp.port(), "POST", "/sessions/term_target/message",
|
||||
"{\"content\":\"config.yaml\",\"turnId\":\"" + turnId + "\"}", null);
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
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"));
|
||||
HttpResponse<String> ownedResponse = owned.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(200, ownedResponse.statusCode(), ownedResponse.body());
|
||||
assertEquals("done", mapper.readTree(ownedResponse.body()).path("reply").asText(null),
|
||||
"the real owner's answer must resolve the worker's turn");
|
||||
} finally {
|
||||
ownerApp.stop();
|
||||
attackerApp.stop();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, but the delegation is a blocking send ({@code wait:true}) instead of a ticket —
|
||||
* the owning caller's own HTTP request is the one that surfaces the worker's question and
|
||||
* later carries the real answer. Covers the REST adapter's ANSWER gate for a blocking-send
|
||||
* delegation.
|
||||
*/
|
||||
@Test
|
||||
void restAnswerIsRefusedForADifferentCallerOnABlockingSendDelegationButTheRealOwnerSucceeds() 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 ownerApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID, Map.of("term_a", "lead-owner"));
|
||||
Javalin attackerApp = startOnSharedService(messages, herdr, 9001L, Map.of("term_shell", "lead-attacker"));
|
||||
try {
|
||||
ObjectMapper mapper = new ObjectMapper();
|
||||
CompletableFuture<HttpResponse<String>> blocking = CompletableFuture.supplyAsync(() -> {
|
||||
try {
|
||||
return send(ownerApp.port(), "POST", "/sessions/term_target/message",
|
||||
"{\"content\":\"long task\",\"wait\":true,\"timeoutMs\":5000}", null);
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_target"), "the blocking send should have opened its rendezvous waiter");
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> messages.ask("term_target", "which config file?", 5000));
|
||||
|
||||
HttpResponse<String> questionResponse = blocking.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(202, questionResponse.statusCode(), questionResponse.body());
|
||||
JsonNode question = mapper.readTree(questionResponse.body());
|
||||
assertEquals("question", question.get("status").asText());
|
||||
String turnId = question.get("turnId").asText();
|
||||
|
||||
HttpResponse<String> hijacked = send(attackerApp.port(), "POST", "/sessions/term_target/message",
|
||||
"{\"content\":\"hijack\",\"turnId\":\"" + turnId + "\"}", null);
|
||||
assertEquals(403, hijacked.statusCode(), hijacked.body());
|
||||
assertTrue(hijacked.body().contains("not_turn_owner"),
|
||||
"a different caller's answer must be refused as not_turn_owner: " + hijacked.body());
|
||||
assertEquals("term_target", rendezvous.askSession(turnId),
|
||||
"a refused answer must leave the ask turn open");
|
||||
|
||||
CompletableFuture<HttpResponse<String>> owned = CompletableFuture.supplyAsync(() -> {
|
||||
try {
|
||||
return send(ownerApp.port(), "POST", "/sessions/term_target/message",
|
||||
"{\"content\":\"config.yaml\",\"turnId\":\"" + turnId + "\"}", null);
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
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"));
|
||||
HttpResponse<String> ownedResponse = owned.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(200, ownedResponse.statusCode(), ownedResponse.body());
|
||||
assertEquals("done", mapper.readTree(ownedResponse.body()).path("reply").asText(null),
|
||||
"the real owner's answer must resolve the worker's turn");
|
||||
} finally {
|
||||
ownerApp.stop();
|
||||
attackerApp.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
|
||||
|
||||
@@ -2436,4 +2436,160 @@ class SessionManagerTest {
|
||||
return "threw:" + e.getClass().getName() + ":" + e.getMessage();
|
||||
}
|
||||
}
|
||||
|
||||
// --- fleetd #702: spawnedMemberRole must still answer for a pane mid-teardown ----------------
|
||||
|
||||
/**
|
||||
* A {@link Worktrees} test double whose {@code hasUncommitted} runs an injected hook before
|
||||
* answering. This is what lets a test resolve a releasing pane's terminal from inside the
|
||||
* window {@link SessionManager#release} opens between removing the registry entry and
|
||||
* actually stopping the pane — the hook runs synchronously on the release call's own thread,
|
||||
* at the exact point {@code release} shells out to {@code git status}, so there is no sleep
|
||||
* and no race to land in it.
|
||||
*/
|
||||
private static final class HookedWorktrees implements Worktrees {
|
||||
private Runnable hook;
|
||||
private boolean dirty = false;
|
||||
|
||||
HookedWorktrees onHasUncommitted(Runnable hook) {
|
||||
this.hook = hook;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
return "/wt/" + branch.replace('/', '_');
|
||||
}
|
||||
|
||||
@Override
|
||||
public void remove(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
if (hook != null) {
|
||||
hook.run();
|
||||
}
|
||||
return dirty;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String repoRoot(String cwd) {
|
||||
return "/repo";
|
||||
}
|
||||
|
||||
@Override
|
||||
public java.util.Optional<String> snapshot(String worktreePath, String branch, String message) {
|
||||
return java.util.Optional.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WipRefStats wipRefs(String repoRoot) {
|
||||
return new WipRefStats(0, 0L);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnedMemberRoleResolvesTheLiveRegisteredRole() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller", null);
|
||||
|
||||
assertEquals(s.role(), sessions.spawnedMemberRole(s.terminalId()),
|
||||
"a registered session resolves to its own role");
|
||||
assertNull(sessions.spawnedMemberRole("term_unknown"),
|
||||
"a terminal with no session at all resolves to null");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnedMemberRoleIsNullOnceReleaseFullyCompletes() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr, new HookedWorktrees());
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-702a", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertNull(sessions.spawnedMemberRole(s.terminalId()),
|
||||
"once release has fully finished, the terminal is neither registered nor releasing");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnedMemberRoleStillAnswersBetweenTheRegistryRemovalAndThePaneStop() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
java.util.concurrent.atomic.AtomicReference<MemberRole> duringWindow = new java.util.concurrent.atomic.AtomicReference<>();
|
||||
HookedWorktrees worktrees = new HookedWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-702b", null));
|
||||
worktrees.onHasUncommitted(() -> {
|
||||
// This runs from INSIDE release()'s git-status shell-out: the registry entry is
|
||||
// already gone, but the pane has not stopped yet — a real call landing inside the
|
||||
// exact window fleetd #702 reports, so no sleep and no race is needed to reach it.
|
||||
assertTrue(sessions.get(s.paneId()).isEmpty(),
|
||||
"sanity: the registry entry is already gone at this point");
|
||||
duringWindow.set(sessions.spawnedMemberRole(s.terminalId()));
|
||||
});
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(s.role(), duringWindow.get(),
|
||||
"spawnedMemberRole must still answer the live role while the pane is mid-teardown, "
|
||||
+ "not only while the session is still in the registry");
|
||||
assertNull(sessions.spawnedMemberRole(s.terminalId()),
|
||||
"and once release has fully finished, the window is closed too");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code Releasing.enter} and {@code Releasing.leave} are called directly here, never through
|
||||
* {@link SessionManager#release} or {@link SessionManager#spawnedMemberRole}, so a test naming
|
||||
* one of them exercises only that one — a regression in the other can never hide behind it.
|
||||
*/
|
||||
@Test
|
||||
void releasingLeaveStepsDownADepthGreaterThanOneInsteadOfRemovingIt() {
|
||||
SessionManager.Releasing depthTwo = new SessionManager.Releasing(2, "term_a", MemberRole.DEV);
|
||||
|
||||
SessionManager.Releasing afterLeave = depthTwo.leave();
|
||||
|
||||
assertNotNull(afterLeave,
|
||||
"depth 2 means another release of the SAME pane is still mid-teardown; leave() must "
|
||||
+ "step the depth down, never remove the marker outright — removing it here "
|
||||
+ "is what a plain Set would do, and would reopen the window the still-in-"
|
||||
+ "flight release is relying on staying closed");
|
||||
assertEquals(1, afterLeave.depth());
|
||||
assertEquals("term_a", afterLeave.terminalId());
|
||||
assertEquals(MemberRole.DEV, afterLeave.role());
|
||||
}
|
||||
|
||||
@Test
|
||||
void releasingEnterPreservesThePriorTerminalWhenTheOverlappingCallHasNoSessionOfItsOwn() {
|
||||
SessionManager.Releasing prior = new SessionManager.Releasing(1, "term_a", MemberRole.DEV);
|
||||
|
||||
SessionManager.Releasing afterEnter = SessionManager.Releasing.enter(prior, null);
|
||||
|
||||
assertEquals(2, afterEnter.depth(), "depth still increments whether or not this enter knows its session");
|
||||
assertEquals("term_a", afterEnter.terminalId(),
|
||||
"an overlapping release that finds the registry entry already gone has no session of "
|
||||
+ "its own to pass as known, and must not blank out the terminal the first "
|
||||
+ "call already recorded — that terminal is what spawnedMemberRole matches "
|
||||
+ "against");
|
||||
assertEquals(MemberRole.DEV, afterEnter.role());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user