Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha bf2d940c26 fleetd #669: fix fleetd.example.yaml's false scan-exclusion claims; cover the shipped shape
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 2m10s
fleetd.example.yaml claimed member spaces are excluded from the lead-tab scan and that
a lead's workspace default is "leads" and must not match a member workspace. Both are
false: production always passes an empty excludedWorkspaceLabels set (CB-558), and a
lead's workspace defaults to the same shared "fleet" space the members use. Rewrite
both paragraphs to name the defenses that actually exist: the startup tabPrefix refusal
and CallerResolver's roster-first precedence.

Add the test nobody wrote: a lead is still discovered when its workspace label equals
the member workspace, with an empty exclusion set. Correct
aTabInAWorkerSpaceIsNeverALeadEvenWhenItsLabelMatches's javadoc, which overclaimed that
the excludedWorkspaceLabels parameter is the live production guard. Drop the stale line
number from FleetdAssemblyLeadTabScannerExclusionTest's assertion message and javadoc.
2026-10-04 06:31:20 +02:00
19 changed files with 202 additions and 2130 deletions
+4 -6
View File
@@ -100,9 +100,7 @@ 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 — 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.
**A worker's ask waits ~55 seconds, and no nudge makes that longer** — so never
that answers it. **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
@@ -170,7 +168,7 @@ you decide.
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — `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 |
| 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 |
| 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}` |
@@ -259,8 +257,8 @@ of those is refused at the gate, not queued.
Two limits worth knowing before you hit them. **You cannot reach a worker** — not even to help one —
because a worker belongs to the lead that spawned it, and routing around that would make you a
second orchestrator with no plan. Send to the lead instead. And **you cannot read a ticket**, so you
cannot collect a delegation's reply: `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.
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.
Being named buys you a channel, not authority. Your `fleet_send` to a lead is coordination between
peers: the lead owes you no obedience, and you owe it none.
@@ -75,9 +75,6 @@ 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,6 +249,17 @@ 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}.
*
@@ -2776,10 +2776,9 @@ public record FleetConfig(
});
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: " + String.join("; ", bad)
+ ". A member labelled that way, while its pane carries no entry in the "
+ "spawned-member roster, is read back as a lead or collaborator and granted "
+ "that identity's authority. Change one of the two so member tabs cannot be "
+ "confused with a lead's or collaborator's tab.");
+ ". Every member labelled that way would be read back as a lead or "
+ "collaborator and granted that identity's authority. Change one of the two "
+ "so member tabs cannot be confused with a lead's or collaborator's tab.");
}
List<String> collisions = new ArrayList<>();
@@ -2878,8 +2877,7 @@ public record FleetConfig(
* the focused tab rather than its own, so it can land inside a lead's or collaborator's own
* labelled tab. {@link dev.ltms.fleet.herdr.LeadTabScanner} identifies a lead or collaborator
* purely by that tab's label — it does not exclude the member space — so a member that ends up
* there, while its pane carries no entry in the spawned-member roster, is read back as that
* lead or collaborator and granted that identity's authority.
* there would be read back as that lead or collaborator and granted that identity's authority.
*
* <p>Only an entry with a non-blank {@code tab} is in scope: one with no {@code tab} feeds
* nothing into {@link dev.ltms.fleet.herdr.LeadTabScanner}, so it creates no hazard here.
@@ -2910,9 +2908,8 @@ public record FleetConfig(
}
throw new IllegalStateException("refusing to start: profile(s) " + bad
+ " use placement: pane while fleet.leaders or fleet.collaborators names a tab. A "
+ "pane-placed member can land inside that labelled tab, and while its pane "
+ "carries no entry in the spawned-member roster, it is read back as the lead or "
+ "collaborator and granted that identity's authority. Set placement: tab for "
+ "pane-placed member can land inside that labelled tab and be read back as the "
+ "lead or collaborator, granted that identity's authority. Set placement: tab for "
+ "each named profile, or remove the tab from every fleet.leaders and "
+ "fleet.collaborators entry.");
}
@@ -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), caller);
return answer(messages, turnId, content, timeoutMs(a));
}
// 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(), caller)
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), caller);
? sendAsync(messages, target, content, onAccepted, workers.profiles())
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
};
// 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"), callerTerminal(exchange));
return status(messages, str(req.arguments(), "sessionId"));
};
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, callerTerminal(exchange));
return poll(messages, leadChannel, str(a, "ticket"), target, coordId);
};
// 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.
@@ -560,10 +560,8 @@ public final class FleetMcp {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
callerTerminal(exchange),
callers.collaborators(), collaboratorsVisibleTo(principal(exchange)),
new CoordinationSource(leadChannel, peers),
coordinatorVisibleTo(principal(exchange)),
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)));
coordinatorVisibleTo(principal(exchange)));
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -745,53 +743,7 @@ public final class FleetMcp {
return caller.isPrimary();
}
/**
* fleetd #703: who may see {@code fleet_list}'s {@code collaborators} array — a roster of the
* human-opened tabs this daemon recognises as named peers. Visible to exactly the roles that
* may {@link Authz.Action#SEND} to a named peer ({@code Authz.java}'s {@code SEND} case):
* the primary, an architect, and a collaborator (reaching another collaborator or a lead). A
* worker holds {@code READ} but can never {@code SEND} to a collaborator, so listing them to a
* worker would expose which human tabs exist on the host with no use to that caller. Split out
* for the same reason as {@link #coordinatorVisibleTo}: 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 collaboratorsVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
}
/**
* 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.
*/
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
private static String callerTerminal(McpSyncServerExchange exchange) {
Object v = exchange.transportContext().get(CALLER_TERMINAL);
String s = v == null ? null : v.toString();
@@ -908,8 +860,7 @@ 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,
String callerTerminal) {
Long timeoutMs, Runnable onAccepted, Set<String> profiles) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -919,7 +870,7 @@ public final class FleetMcp {
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
try {
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerTerminal), timeout);
return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout);
} catch (HerdrException e) {
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
}
@@ -929,15 +880,13 @@ 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,
String callerTerminal) {
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs) {
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, callerTerminal), timeout);
return formatReply(messages.answer(turnId, content, timeout), timeout);
}
/**
@@ -986,8 +935,6 @@ 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,
@@ -1009,17 +956,6 @@ 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");
}
@@ -1027,7 +963,7 @@ public final class FleetMcp {
if (targetError != null) {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
String ticket = messages.sendAsync(sessionId, content, onAccepted);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
@@ -1206,24 +1142,9 @@ 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);
}
@@ -1237,7 +1158,7 @@ public final class FleetMcp {
if (isBlank(ticket)) {
return error("ticket (or target) is required");
}
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
MessageService.TaskView v = messages.poll(ticket);
if (v == null) {
return error("unknown ticket: " + ticket + " (never issued, or expired)");
}
@@ -1344,20 +1265,16 @@ 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} — 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.
* 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.
*/
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerTerminal) {
static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
if (isBlank(sessionId)) {
return error("sessionId is required");
}
try {
String base = messages.status(sessionId).name().toLowerCase();
MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerTerminal);
MessageService.PendingAsk ask = messages.pendingAsk(sessionId);
if (ask == null) {
return text(base);
}
@@ -1844,8 +1761,7 @@ public final class FleetMcp {
Map<String, String> leads, String selfTerm,
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, true, true);
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1854,8 +1770,7 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
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, true, true);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, CoordinationSource.none(), false);
}
/**
@@ -1872,8 +1787,7 @@ public final class FleetMcp {
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
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, true, true);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1882,8 +1796,7 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
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, true, true);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1905,8 +1818,7 @@ public final class FleetMcp {
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
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, true, true);
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1923,23 +1835,14 @@ 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,
boolean leadsVisible, boolean membersVisible) {
CoordinationSource coordination, boolean callerIsPrimary) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, callerIsPrimary, leadsVisible, membersVisible);
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, callerIsPrimary);
}
/**
@@ -1949,17 +1852,6 @@ public final class FleetMcp {
* matter); the one caller that matters for caching, {@code fleet_list}'s MCP handler, passes
* its own single long-lived instance instead (see {@code FleetMcp}'s {@code leadContextGauge}
* field).
*
* @param collaborators terminal_id → collaborator name, live from the resolver
* ({@link dev.ltms.fleet.auth.CallerResolver#collaborators()})
* @param collaboratorsVisible whether this caller may see the {@code collaborators} array (see
* {@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,
@@ -1968,56 +1860,32 @@ public final class FleetMcp {
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
Map<String, String> collaborators, boolean collaboratorsVisible,
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
CoordinationSource coordination, boolean callerIsPrimary) {
try {
// 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();
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();
// 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<>();
// 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("leads", leadRows); result.put("members", out);
result.put("healthCoverage", healthCoverage.value().get());
result.put("loopHealth", Map.of(
"statusPoller", loopHealth.statusPoller().get().name(),
"sessionReaper", loopHealth.sessionReaper().get().name()));
// fleetd #703: a collaborator tab is a person's own tab, so the row is assembled and
// included only for the roles that may SEND to a named peer -- gate BEFORE assembling
// it, same reason as the coordinator row just below: the key must be absent for a
// worker, never present-and-empty.
if (collaboratorsVisible) {
result.put("collaborators", collaborators.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> collaboratorRow(e.getKey(), e.getValue()))
.toList());
}
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
// key is absent rather than present-and-empty.
@@ -2329,20 +2197,6 @@ public final class FleetMcp {
return m;
}
/**
* fleetd #703: one row of the {@code collaborators} array — a named peer's registry name and
* the herdr {@code terminal_id} a {@code fleet_send} must target to reach it. No status, no
* context gauge, no profile: a collaborator is never spawned and carries no profile, so those
* fields have no meaning for it, and the scan behind {@code terminal} only reports a tab it
* actually found in herdr, so a tab nobody has open does not appear here at all.
*/
private static Map<String, Object> collaboratorRow(String terminal, String name) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("name", name);
m.put("sessionId", terminal);
return m;
}
/**
* Reads the lead context gauge for one lead. {@code configDir} is this lead's configured
* {@code CLAUDE_CONFIG_DIR} override (see {@link LeadConfigDirSource}), derived from
@@ -2537,14 +2391,7 @@ public final class FleetMcp {
private static McpSchema.Tool listTool() {
return tool(FleetTool.LIST.wireName(),
"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 "
"List the whole fleet the bridge tracks, in two parts. '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 "
@@ -2575,12 +2422,7 @@ public final class FleetMcp {
+ "msgId/from/preview, never the full body), and one row per coordinator.peers "
+ "coord-id ('peers': coordId/reachable, plus pending/consumers when reachable) — "
+ "this is peer DISCOVERY for cross-host leads, distinct from the local 'leads' "
+ "array above. It is omitted entirely when no coordinator is configured. A "
+ "'collaborators' array, visible only to the primary, an architect, and a "
+ "collaborator (never a worker), reports every other named peer tab this daemon "
+ "recognises: each row is 'name' (its registry name) and 'sessionId' (the "
+ "terminal id a fleet_send targets to reach it). Absent entirely for a worker, "
+ "whatever is configured.",
+ "array above. It is omitted entirely when no coordinator is configured.",
objectSchema(Map.of(), List.of()));
}
@@ -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, String)} and the turn resumes.
* {@link #answer(String, String, long)} and the turn resumes.
*/
QUESTION,
/** Timed out after the message was delivered — the worker is still working. */
@@ -120,16 +120,10 @@ public final class MessageService {
/** Another send to this session was in flight for the whole window. */
BUSY,
/**
* 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.
* 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.
*/
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
STALE_TURN
}
/**
@@ -139,7 +133,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, String)}), else {@code null}
* {@link #answer(String, String, long)}), else {@code null}
*/
public record Reply(Outcome outcome, String text, String turnId) {
/** A reply with no correlation id (the common terminal outcomes). */
@@ -281,19 +275,10 @@ 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, String creatorTerminal) {
private Task(String ticket, String target, LongSupplier nowNanos) {
this.ticket = ticket;
this.target = target;
this.creatorTerminal = creatorTerminal;
this.createdNanos = nowNanos.getAsLong();
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
@@ -671,7 +656,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, NOT_TURN_OWNER -> null; // not a completed delegation
case STALE_TURN, QUESTION -> null; // not a completed delegation
};
}
@@ -930,17 +915,14 @@ 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. {@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.
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses.
*/
public Reply send(String target, String content, long timeoutMillis, String callerTerminal) {
return send(target, content, timeoutMillis, null, callerTerminal);
public Reply send(String target, String content, long timeoutMillis) {
return send(target, content, timeoutMillis, null);
}
/**
* As {@link #send(String, String, long, String)}, but with an accepted-delivery hook.
* As {@link #send(String, String, long)}, 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
@@ -951,13 +933,12 @@ 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, String callerTerminal) {
return send(target, content, timeoutMillis, onAccepted, null, callerTerminal);
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted) {
return send(target, content, timeoutMillis, onAccepted, null);
}
/** 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,
String callerTerminal) {
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task) {
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
@@ -977,7 +958,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, Rendezvous.Owner.of(callerTerminal));
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
// 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.
@@ -1156,30 +1137,19 @@ 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, String callerTerminal) {
public Reply answer(String turnId, String content, long timeoutMillis) {
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, owner);
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
// 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
@@ -1309,20 +1279,8 @@ 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, creatorTerminal);
Task task = new Task(ticket, target, nowNanos);
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
@@ -1353,7 +1311,7 @@ public final class MessageService {
}
asyncExecutor.submit(() -> {
try {
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorTerminal);
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
if (result.outcome() == Outcome.QUESTION) {
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
// just after resolveQuestion wakes this thread.
@@ -1385,31 +1343,15 @@ public final class MessageService {
}
/**
* 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
* 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.
*/
public TaskView poll(String ticket, String callerTerminal) {
public TaskView poll(String ticket) {
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;
@@ -1445,17 +1387,6 @@ 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.
@@ -1775,17 +1706,16 @@ public final class MessageService {
}
/**
* 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)}).
* 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}).
*/
public PendingAsk pendingAsk(String workerSession, String callerTerminal) {
public PendingAsk pendingAsk(String workerSession) {
for (Task task : tasks.values()) {
Reply q = task.question;
if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerTerminal)) {
if (q != null && workerSession.equals(task.target)) {
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, Owner owner) {
private record AskWaiter(String session, CompletableFuture<String> answer) {
}
/**
@@ -70,33 +70,7 @@ public final class Rendezvous {
public record AskTicket(String turnId, CompletableFuture<String> answer, boolean fresh) {
}
/**
* 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<>();
private final ConcurrentHashMap<String, CompletableFuture<Resolution>> 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<>();
@@ -115,17 +89,8 @@ 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<>();
ForwardWaiter existing = waiters.putIfAbsent(session, new ForwardWaiter(owner, waiter));
CompletableFuture<Resolution> existing = waiters.putIfAbsent(session, waiter);
if (existing != null) {
throw new IllegalStateException(
"rendezvous double-open for session " + session + " — a waiter is already registered");
@@ -140,13 +105,7 @@ 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.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();
waiters.remove(session, waiter);
}
/** Whether a send is currently awaiting a resolution for {@code session}. */
@@ -160,8 +119,7 @@ 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) {
ForwardWaiter w = waiters.get(session);
return w == null ? null : w.future();
return waiters.get(session);
}
/**
@@ -189,7 +147,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, ownerOf(session));
AskWaiter waiter = new AskWaiter(session, answer);
asks.put(newTurnId, waiter);
minted[0] = waiter;
return newTurnId;
@@ -225,16 +183,6 @@ 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.
@@ -297,7 +245,7 @@ public final class Rendezvous {
}
private boolean complete(String session, Resolution resolution) {
ForwardWaiter waiter = waiters.get(session);
return waiter != null && waiter.future().complete(resolution);
CompletableFuture<Resolution> waiter = waiters.get(session);
return waiter != null && waiter.complete(resolution);
}
}
@@ -663,8 +663,6 @@ 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());
@@ -690,19 +688,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, callerTerminal), timeout);
writeReply(ctx, id, messages.answer(turnId, content, timeout), 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, null, callerTerminal);
String ticket = messages.sendAsync(id, content);
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
return;
}
try {
writeReply(ctx, id, messages.send(id, content, timeout, callerTerminal), timeout);
writeReply(ctx, id, messages.send(id, content, timeout), timeout);
} catch (HerdrException e) {
herdrError(ctx, e);
}
@@ -721,10 +719,6 @@ 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).
@@ -739,10 +733,9 @@ 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, 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.
// 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.
"status", switch (reply.outcome()) {
case TIMED_OUT_WORKING -> "working";
case TIMED_OUT_QUEUED -> "queued";
@@ -752,7 +745,7 @@ public final class FleetApp {
case BUSY -> "busy";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN, NOT_TURN_OWNER -> "done"; // unreachable
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable
},
"detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED
? "no reply within " + timeout + "ms; delivery is unconfirmed — the "
@@ -881,12 +874,10 @@ public final class FleetApp {
body.put("sessionId", id);
body.put("status", messages.status(id).name().toLowerCase());
body.put("ready", deliverable.test(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());
// 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);
if (ask != null) {
body.put("question", ask.question());
body.put("turnId", ask.turnId());
@@ -903,8 +894,7 @@ public final class FleetApp {
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
return;
}
Principal caller = ctx.attribute(CALLER);
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.terminal());
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"));
if (v == null) {
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
return;
@@ -30,22 +30,6 @@ 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}) {
@@ -460,6 +460,7 @@ 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
@@ -850,8 +850,7 @@ class FleetConfigTest {
/**
* The hazard this guard closes: a pane-placed member lands inside the focused tab rather than
* its own, so it can land inside a lead's labelled tab and, while its pane carries no entry in
* the spawned-member roster, be read back as that lead.
* its own, so it can land inside a lead's labelled tab and be read back as that lead.
*/
@Test
void aPanePlacedProfileWithALeadTabRefusesToStart(@TempDir Path dir) throws Exception {
@@ -1088,9 +1087,8 @@ class FleetConfigTest {
}
/**
* A member tabLabel that can render as a configured collaborator tab is the same hazard as the
* lead case above — while its pane carries no entry in the spawned-member roster, a member
* labelled that way is read back as the collaborator.
* fleetd #669: a member tabLabel that can render as a configured collaborator tab is the same
* hazard as the lead case above — a member labelled that way is read back as the collaborator.
*/
@Test
void aProfileTabLabelOverrideMatchingACollaboratorTabRefusesToStart(@TempDir Path dir)
@@ -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*[,)]",
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*\\)\\s*;",
Pattern.DOTALL);
Matcher m = trailingArg.matcher(handlerBlock);
assertTrue(m.find(), "could not locate listFleet(...)'s coordinator boolean argument in the "
assertTrue(m.find(), "could not locate listFleet(...)'s trailing 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,
@@ -366,197 +366,6 @@ class FleetMcpAuthzTest {
+ "calling, not pass a literal boolean -- found: " + trailing);
}
// --- fleetd #703: who may see fleet_list's collaborators array -------------------------------
/**
* fleetd #703: {@link FleetMcp#collaboratorsVisibleTo} is the whole policy decision for
* {@code fleet_list}'s {@code collaborators} array. Visible to exactly the roles that may
* {@code SEND} to a named peer -- the primary, an architect, and a collaborator itself -- never
* a worker, which holds {@code READ} but can never {@code SEND} to a collaborator, and never an
* anonymous caller.
*/
@Test
void onlyPrimaryArchitectAndCollaboratorMaySeeTheCollaboratorsArray() {
assertTrue(FleetMcp.collaboratorsVisibleTo(PRIMARY), "the primary must see the collaborators array");
assertTrue(FleetMcp.collaboratorsVisibleTo(ARCH_DESIGN), "an architect must see the collaborators array");
assertTrue(FleetMcp.collaboratorsVisibleTo(COLLABORATOR), "a collaborator must see its own peer roster");
assertFalse(FleetMcp.collaboratorsVisibleTo(WORKER_A),
"a worker holds READ but can never SEND to a collaborator, so it must not see the array");
assertFalse(FleetMcp.collaboratorsVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* fleetd #703, 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 both threads {@code callers.collaborators()} into the payload and
* asks {@code collaboratorsVisibleTo(principal(exchange))} for the visibility flag, rather than a
* literal boolean or an empty map.
*/
@Test
void theFleetListHandlerActuallyConsultsCollaboratorsVisibleTo() 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 assertions 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("callers.collaborators()"),
"the fleet_list handler must thread callers.collaborators() into listFleet(...), not an "
+ "empty or literal map -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("collaboratorsVisibleTo(principal(exchange))"),
"the fleet_list handler must ask collaboratorsVisibleTo(principal(exchange)) who is "
+ "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) ------------------------------------
/**
@@ -121,7 +121,6 @@ class FleetMcpLeadContextGaugeWiringTest {
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
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, true, true);
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, FleetMcp.CoordinationSource.none(), false);
}
}
@@ -10,7 +10,6 @@ import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
@@ -77,7 +76,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, null));
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -103,7 +102,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(), null));
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of()));
// 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 +180,7 @@ class FleetMcpTest {
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
@@ -279,7 +278,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, null));
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -288,89 +287,6 @@ 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);
@@ -380,7 +296,7 @@ class FleetMcpTest {
@Test
void sendTimesOutWithAWorkingNote() {
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null);
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of());
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
}
@@ -396,7 +312,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(), null));
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of()));
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -416,13 +332,13 @@ class FleetMcpTest {
@Test
void sendRejectsMissingArgs() {
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null).isError());
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null).isError());
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError());
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of()).isError());
}
@Test
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null);
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"));
McpSchema.CallToolResult async = FleetMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
assertTrue(blocking.isError());
@@ -441,7 +357,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, null);
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles);
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
}
@@ -533,7 +449,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(), null));
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of()));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -555,7 +471,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, null));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L));
// The worker's ask returns the answer — it resumes the same turn.
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
@@ -581,7 +497,7 @@ class FleetMcpTest {
@Test
void answerToAStaleTurnIsAnError() {
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null);
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L);
assertTrue(res.isError());
assertTrue(textOf(res).contains("no longer open"), textOf(res));
}
@@ -775,7 +691,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, true, true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
String out = textOf(res);
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
@@ -804,7 +720,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, true, true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
String out = textOf(res);
assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out);
@@ -863,82 +779,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));
Principal worker = Principal.worker("term_a", 1);
boolean callerIsPrimary = Principal.worker("term_a", 1).isPrimary();
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()), worker.isPrimary(),
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
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);
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("\"leads\""), "the rest of the result must still be present: " + out);
assertTrue(out.contains("\"members\""), out);
assertTrue(out.contains("\"healthCoverage\""), out);
}
@@ -954,19 +807,16 @@ 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));
Principal architect = Principal.architect("lead-designer", "term_design", 400);
boolean callerIsPrimary = Principal.architect("lead-designer", "term_design", 400).isPrimary();
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()), architect.isPrimary(),
FleetMcp.leadsVisibleTo(architect), FleetMcp.membersVisibleTo(architect));
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
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);
}
/**
@@ -991,7 +841,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, true, true));
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true));
assertTrue(gatedAsPrimary.contains("\"coordinator\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"selfId\":\"mac-opus\""), gatedAsPrimary);
@@ -1000,8 +850,6 @@ 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
@@ -1020,7 +868,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, true, true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
String out = textOf(res);
assertTrue(out.contains("\"msgId\":\"m1\""), out);
@@ -1031,83 +879,6 @@ class FleetMcpTest {
assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out);
}
// --- fleetd #703: fleet_list's collaborators array -------------------------------------------
/** Calls the canonical {@code listFleet} overload directly, so a test can set the collaborators
* payload and its visibility independently of a real {@code Principal} / MCP exchange. */
private static McpSchema.CallToolResult listFleetWithCollaborators(FakeHerdr h,
Map<String, String> collaborators, boolean collaboratorsVisible) {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
return FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
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, true, true);
}
/**
* fleetd #703 acceptance A: a visible caller with one configured collaborator gets a
* {@code collaborators} row whose key ({@code CallerResolver.collaborators()}'s
* {@code terminal_id -> name} entry) lands as that row's {@code sessionId}, and {@code leads}/
* {@code members} are unaffected by the new key.
*/
@Test
void listReportsACollaboratorsRowKeyedByTheCollaboratorsTerminalId() {
FakeHerdr h = new FakeHerdr();
String out = textOf(listFleetWithCollaborators(h, Map.of("term_collab", "ops"), true));
assertTrue(out.contains("\"collaborators\":["), out);
assertTrue(out.contains("\"name\":\"ops\""), out);
assertTrue(out.contains("\"sessionId\":\"term_collab\""), out);
assertTrue(out.contains("\"leads\":[]"), out);
assertTrue(out.contains("\"members\":[]"), out);
}
/**
* fleetd #703 acceptance A, control half: with no collaborator configured, the key is absent
* (caller cannot see it) or an empty array (caller can), and {@code leads}/{@code members} are
* unchanged either way.
*/
@Test
void listOmitsOrEmptiesCollaboratorsWhenNoneAreConfigured() {
FakeHerdr h = new FakeHerdr();
String visibleButEmpty = textOf(listFleetWithCollaborators(h, Map.of(), true));
assertTrue(visibleButEmpty.contains("\"collaborators\":[]"), visibleButEmpty);
assertTrue(visibleButEmpty.contains("\"leads\":[]"), visibleButEmpty);
assertTrue(visibleButEmpty.contains("\"members\":[]"), visibleButEmpty);
String notVisible = textOf(listFleetWithCollaborators(h, Map.of(), false));
assertFalse(notVisible.contains("\"collaborators\""), notVisible);
assertTrue(notVisible.contains("\"leads\":[]"), notVisible);
assertTrue(notVisible.contains("\"members\":[]"), notVisible);
}
/**
* fleetd #703 acceptance B: a worker must not see the {@code collaborators} array at all, while
* an architect -- a role that also holds READ, same as a worker -- does see it. Both halves are
* asserted: a test that only checked the worker-hidden half would pass even if the feature were
* never wired up for anyone.
*/
@Test
void listOmitsCollaboratorsForAWorkerAndIncludesThemForAnArchitect() {
FakeHerdr h = new FakeHerdr();
Map<String, String> collaborators = Map.of("term_collab", "ops");
String asWorker = textOf(listFleetWithCollaborators(h, collaborators,
FleetMcp.collaboratorsVisibleTo(Principal.worker("term_w", 1))));
assertFalse(asWorker.contains("\"collaborators\""),
"a worker must not see the collaborators array: " + asWorker);
String asArchitect = textOf(listFleetWithCollaborators(h, collaborators,
FleetMcp.collaboratorsVisibleTo(Principal.architect("design", "term_arch", 2))));
assertTrue(asArchitect.contains("\"collaborators\":["),
"an architect must see the collaborators array: " + asArchitect);
assertTrue(asArchitect.contains("\"sessionId\":\"term_collab\""), asArchitect);
}
/**
* fleetd #421: {@code mailbox.pending} counts only broker-ready messages, so a blocked lead's
* normal, healthy state is {@code "pending": 0} next to a non-empty {@code held[]} — which
@@ -1129,7 +900,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, true, true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
String out = textOf(res);
assertTrue(out.contains("\"pending\":0"), out);
@@ -1158,7 +929,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, true, true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
String out = textOf(res);
assertTrue(out.contains("\"heldDurable\":false"),
@@ -1227,7 +998,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, true, true);
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true);
String out = textOf(res);
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out);
@@ -1943,14 +1714,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", null);
new MessageService(blockedAgents, new Injector(blockedAgents), rendezvous), "term_a");
assertNotEquals(Boolean.TRUE, res.isError());
assertEquals("blocked", textOf(res));
}
/**
* 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
* 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
* window it opened with is far shorter than that cadence.
*/
@Test
@@ -1973,7 +1744,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, null);
McpSchema.CallToolResult res = FleetMcp.status(messages, T);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.startsWith("idle"), "the live status must still lead the text: " + out);
@@ -1984,64 +1755,7 @@ 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, 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"));
() -> messages.answer(turnId, "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -1,280 +0,0 @@
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, null));
return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis));
}
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, (String) null));
CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000));
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, null));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
// 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, null));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
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, null));
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500));
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, null)));
lateAnswer.complete(messages.answer(turnId, "too late", 500)));
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, null);
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
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, null);
MessageService.Reply r = messages.answer(turnId, "too late", 500);
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, null);
MessageService.Reply r = messages.send(T, "never delivered", 50);
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, null));
CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 300));
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, null);
MessageService.Reply reply = messages.send(T, "race delivery", 50);
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, null));
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150));
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, null);
MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200);
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, null));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
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, null);
MessageService.Reply probe = messages.send(T, "probe after normal reply", 300);
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, null));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200));
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, null);
MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300);
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, null);
messages.answer(turnId, "config.yaml", 5000);
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, null);
MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300);
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, null);
messages.answer(turnId, "config.yaml", 5000);
caught.set(new AssertionError("expected answer() to throw"));
} catch (Throwable t) {
caught.set(t);
@@ -859,88 +859,17 @@ class MessageServiceTest {
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after interruption", 300, null);
MessageService.Reply probe = messages.send(T, "probe after interruption", 300);
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");
@@ -967,11 +896,11 @@ class MessageServiceTest {
@Test
void concurrentSendToSameSessionWhileFirstHoldsItIsBusy() throws Exception {
CompletableFuture<MessageService.Reply> first =
CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000, null));
CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000));
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, null);
MessageService.Reply busy = messages.send(T, "second", 100);
assertEquals(MessageService.Outcome.BUSY, busy.outcome(),
"a second send while another holds the session is busy, not a hang");
assertNull(busy.text());
@@ -1001,13 +930,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), null));
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L)));
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), null);
MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A));
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");
@@ -1028,7 +957,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), null));
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L)));
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
@@ -1039,7 +968,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), null));
() -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A)));
awaitWaiting();
assertEquals(LEAD_A, reg.nudgeTargetFor(T).orElseThrow(),
"an accepted send after the owner finished becomes the new delegator");
@@ -1060,7 +989,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), null));
() -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L)));
awaitWaiting();
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "L owns the delegation");
@@ -1075,7 +1004,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, null));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // the answering send reopened its forward waiter
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
@@ -1095,7 +1024,7 @@ class MessageServiceTest {
void aThrowingAcceptedHookLeavesNoStaleWaiterOrQueuedOrphan() {
assertThrows(IllegalStateException.class,
() -> messages.send(T, "doomed", 500,
() -> { throw new IllegalStateException("ownership hook failed"); }, null),
() -> { throw new IllegalStateException("ownership hook failed"); }),
"a throwing ownership hook fails the send loudly");
assertFalse(rendezvous.isWaiting(T), "the failed send must not leave a stale rendezvous waiter");
@@ -1223,7 +1152,7 @@ class MessageServiceTest {
@Test
void abandonFailsASendThatIsStillWaitingOnAReleasedSession() throws Exception {
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000, null));
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
awaitWaiting();
assertTrue(messages.abandon(T, "session released"), "a live waiter is abandoned");
@@ -1243,7 +1172,7 @@ class MessageServiceTest {
@Test
void abandonDoesNotOverwriteAnAlreadyResolvedSend() throws Exception {
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000, null));
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
awaitWaiting();
assertTrue(rendezvous.resolve(T, "the real answer"));
@@ -1374,7 +1303,7 @@ class MessageServiceTest {
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, null));
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1414,7 +1343,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, null).outcome());
messages.answer(asking.turnId(), "config.yaml", 200).outcome());
}
@Test
@@ -1451,7 +1380,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, null);
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
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");
@@ -1480,7 +1409,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, null);
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
@@ -1534,7 +1463,7 @@ class MessageServiceTest {
messages.setFinishAsyncTaskRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000, null));
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
@@ -1640,7 +1569,7 @@ class MessageServiceTest {
String turnId = asking.turnId();
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000, null));
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
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");
@@ -1687,7 +1616,7 @@ class MessageServiceTest {
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150, null);
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150);
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 "
@@ -1753,7 +1682,7 @@ class MessageServiceTest {
MessageService.TaskView asking = messages.poll(first);
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, null));
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1787,7 +1716,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, null));
() -> messages.answer(asking1.turnId(), "a1", 5000));
assertEquals("a1", ask1.get(5, TimeUnit.SECONDS).answer());
// Still in the SAME resumed turn — before replying — the worker asks again.
@@ -1808,7 +1737,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, null));
() -> messages.answer(turnId2, "a2", 5000));
assertEquals("a2", ask2.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1875,20 +1804,20 @@ class MessageServiceTest {
}
}
// --- fleet_status pendingAsk() -----------------------------------------------------------
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
@Test
void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception {
assertNull(messages.pendingAsk(T, null), "no async ticket at all -> no pending ask");
assertNull(messages.pendingAsk(T), "no async ticket at all -> no pending ask");
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
assertNull(messages.pendingAsk(T, null), "a plain pending delegation is not a question");
assertNull(messages.pendingAsk(T), "a plain pending delegation is not a question");
injectDelivery();
assertTrue(rendezvous.resolve(T, "done"));
awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertNull(messages.pendingAsk(T, null), "a finished ticket carries no open question either");
assertNull(messages.pendingAsk(T), "a finished ticket carries no open question either");
}
@Test
@@ -1901,151 +1830,21 @@ 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, null);
MessageService.PendingAsk pending = messages.pendingAsk(T);
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, null));
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertNull(messages.pendingAsk(T, null), "an answered question is no longer pending");
assertNull(messages.pendingAsk(T), "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
@@ -2302,7 +2101,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, null));
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -2328,7 +2127,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, null));
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
@@ -2721,7 +2520,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, null);
MessageService.Reply r = messages.send(T, "never delivered", 50);
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
assertTrue(messages.hasQueuedDelivery(T),
@@ -2733,7 +2532,7 @@ class MessageServiceTest {
@Test
void hasQueuedDeliveryClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50, null).outcome());
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
assertTrue(messages.hasQueuedDelivery(T));
// A fresh send accepts delivery (opens its own waiter) — the stale queued fact is cleared.
@@ -2750,7 +2549,7 @@ class MessageServiceTest {
@Test
void hasQueuedDeliveryClearsOnAbandon() {
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50, null).outcome());
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
assertTrue(messages.hasQueuedDelivery(T));
messages.abandon(T, "session released");
@@ -1,187 +0,0 @@
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,80 +132,4 @@ 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,7 +2,6 @@ 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;
@@ -40,8 +39,6 @@ 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;
@@ -142,405 +139,6 @@ 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