Compare commits

...

9 Commits

Author SHA1 Message Date
Dai Ha 7cf6075b79 Correct the configDir note: the count was wrong, the shared dir is the defect
CI / shell-tests (push) Failing after 11s
CI / contract (push) Successful in 44s
CI / build (push) Failing after 1m57s
The previous commit compared 4 configDir lines against 8 total profiles, but
four of those are opencode and never read CLAUDE_CONFIG_DIR. Every claude-code
profile does set one, so the original claim was right and this file said
otherwise.

The real defect is narrower: opus and sonnet name the operator's own config dir,
so for those members the store is shared, and ClaudeCodeLauncher's javadoc
already records that fleetd and the operator's session write that same file.
2026-10-04 10:35:33 +02:00
Dai Ha fbdcd709c9 The plugin addendum claimed every profile sets configDir; it does not
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 42s
CI / build (push) Failing after 2m1s
GET /profiles reports 8 live profiles and fleetd.yaml carries 4 configDir
lines, two of them pointing at the operator's own instance dir. The bullet
stated the blanket claim as a structural limit, so it would have been believed.
The conclusion it supported is unchanged: member-facing assets travel in the
worktree.
2026-10-04 10:31:07 +02:00
Dai Ha 8d3f10d291 Releasing is reached directly by its test, and its javadoc names the real case
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 59s
CI / build (push) Failing after 2m0s
The javadoc said a losing releaseIfCurrent CAS is the call that arrives with no
known session. releaseIfCurrent is only called by the reaper, always with a
non-null expected, so it always has one; the null-known call is an overlapping
release that finds the registry entry already gone. The same wrong claim was in
a test's failure message.

Releasing, enter and leave drop private so SessionManagerTest binds them at
compile time. The six reflection helpers are gone, and a rename now breaks the
build instead of a test run.
2026-10-04 10:25:17 +02:00
Dai Ha 38f4fd64ee Merge PR #724: fleetd #702 — a pane mid-teardown keeps its member identity 2026-10-04 10:25:08 +02:00
Dai Ha 3fc39b981d Only the delegation's creator can answer its member's question
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 56s
CI / build (push) Failing after 1m39s
fleetd #715 gates fleet_send{turnId} on the caller that created the delegation,
so the shipped block had to say so: the step-5 note already covered who can see
a pending question, not who can answer it.
2026-10-04 10:15:24 +02:00
Dai Ha 70328ca0f8 Merge PR #725: fleetd #715 — gate fleet_send{turnId} on the caller that created the delegation
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 1m47s
2026-10-04 10:12:38 +02:00
Dai Ha 2374de28e4 fleetd #715: pin that no production caller uses Rendezvous.open(String)
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Failing after 2m4s
Add a source-scrape test mirroring #718's MessageServicePollUsageTest
shape, for the same residual: a convenience overload that defaults the
turn's owner to null, left in place because deleting it would break 81
test-only call sites across 7 unrelated files.

The scanner is exercised against a file known to hold many real
one-argument rendezvous.open( calls before it is ever pointed at
production, using the identical matching logic for both. A pattern that
cannot find the known calls would also find none in production, and
that is exactly the failure mode a -based git grep regex hit earlier
on this ticket: git grep's -E engine does not treat \b as a word
boundary, so that pattern silently matched nothing anywhere, in clean
code and in the 81 real calls alike.
2026-10-04 09:56:05 +02:00
Dai Ha dd18bd1f38 fleetd #715: gate answer() on the caller that owns the turn
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Failing after 2m15s
Record the turn's owner on the forward rendezvous waiter (Rendezvous.Owner,
a three-state record: no record / unnamed primary / named terminal). A
fresh fleet_ask copies that owner onto the ask turn; a coalesced duplicate
ask keeps the first owner. answer() compares the answering caller against
the stored owner before taking the session lock or reopening the resumed
waiter, and a mismatch returns the new NOT_TURN_OWNER outcome instead of
STALE_TURN, with no rendezvous/task/question cleanup.

send() and answer() both drop their no-caller overloads; every call site
in MessageService, FleetMcp and FleetApp now threads an explicit caller
terminal through. AuthzTest pins that the unnamed-primary null allowance
is safe only because ANONYMOUS never reaches ANSWER.

Covers all four combinations of capture input x adapter: blocking-send and
async-send delegations, answered through both FleetMcp and FleetApp, each
with a hijack attempt refused and the real owner's answer succeeding as a
control.
2026-10-04 09:46:20 +02:00
Dai Ha ed4f4b08ad callerTerminal's javadoc names which callers carry a terminal
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m59s
It said "null if the primary", which reads as every primary. Only the unnamed
primary has no terminal; a named lead carries one, so a reader deciding what a
null means was given the wrong rule.
2026-10-04 09:11:08 +02:00
13 changed files with 813 additions and 173 deletions
+19 -5
View File
@@ -102,6 +102,8 @@ below are the procedure — run them in order, every task, not only the big ones
`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.
**Only that same creator can answer it.** A `turnId` you came by any other way is refused, so an
architect's worker waits for that architect and not for you.
**A worker's ask waits ~55 seconds, and no nudge makes that longer** — so never
brief a worker to "ask me". Decide before you delegate, or give it an explicit default.
**A correction cannot reach a busy member.** A `fleet_send` to a working member is *accepted* and
@@ -167,7 +169,7 @@ you decide.
| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) + `loopHealth` (`RUNNING`, `STALLED`, or `STOPPED` for `statusPoller` and `sessionReaper`) · one peer's state: `fleet_status{sessionId}` |
| Delegate (blocking) | `fleet_send{sessionId, content}` |
| Delegate (long task) | `fleet_send{sessionId, content, wait:false}` → ticket → `fleet_poll{ticket}` |
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId`, and only the caller that created that delegation |
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — `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 |
@@ -332,10 +334,22 @@ must obey belongs in the charter, not here.
(#362). **Read `plugin/` before designing anything about onboarding a project.** Two limits are
structural, not bugs: a plugin cannot carry the role agent files, because
`ClaudeCodeLauncher.java:371` requires `<cwd>/.claude/agents/<role>.md` in the member's own
worktree; and a plugin cannot deliver anything to members at all, because
`ClaudeCodeLauncher.java:285` exports `CLAUDE_CONFIG_DIR` and every Claude profile here sets it,
so a member never reads the operator's plugin store. **The plugin is the lead-side surface;
member-facing assets travel in the worktree.**
worktree; and a plugin reaches a member only through `CLAUDE_CONFIG_DIR`, which
`ClaudeCodeLauncher.java:286` exports with `putIfPresent` — so only for a profile that sets
`configDir`. Every `claude-code` profile does set one (the four without are `opencode`, which
never reads that variable). **But measured 2026-10-04: two of them point at
`~/.ccs/instances/ltms`, which is the operator's own `CLAUDE_CONFIG_DIR` on this host.** So for an
`opus` or `sonnet` member, "a member never reads the operator's plugin store" is false — it reads
the same store, because that store is the one its `configDir` names. It stays true for `local` and
`local-direct`, which point at `~/.ccs/instances/gx10`. `ClaudeCodeLauncher`'s own javadoc names
the related hazard: that file is rewritten on every spawn, so for those two profiles fleetd and the
operator's live session write the same `.claude.json`, and its compare-and-swap "narrows the
lost-update window, it does not close it". Re-measure which profiles share the operator's dir with
`awk '/^profiles:/{i=1;next} /^[a-z]/{i=0} i&&/^ [a-z-]+:$/{p=$1} i&&/configDir:/{print p,$2}'
fleetd/fleetd.yaml` against `echo $CLAUDE_CONFIG_DIR`; delete this note once no profile names the
operator's dir. **Treat the plugin as the lead-side surface and put member-facing assets in the
worktree** — that conclusion holds either way, because a worktree asset does not depend on which
config dir a member reads.
- **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/`
(a submodule with its own remote).
- **A provisioned worktree neutralizes `.mcp.json`, `opencode.json` and `.autoenv`** — the repo's
@@ -480,7 +480,7 @@ public final class FleetMcp {
// Answering a worker's fleet_ask (CB-205): resolve its blocked question and
// block for the worker's reply as it resumes the same turn. This is the same
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
return answer(messages, turnId, content, timeoutMs(a));
return answer(messages, turnId, content, timeoutMs(a), caller);
}
// CB-548: delegator ownership (which lead's reply nudge this worker routes to,
// CB-532) is recorded only once the send is ACCEPTED — MessageService has won the
@@ -491,7 +491,7 @@ public final class FleetMcp {
// 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());
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), caller);
};
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
// is "is this caller a worker at all", and it can only ever reply as itself.
@@ -787,7 +787,11 @@ public final class FleetMcp {
return caller.isPrimary() || caller.isArchitect();
}
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
/**
* The terminal of the caller on this call's connection, or {@code null} when that caller carries
* no terminal, which is only the unnamed primary. A named lead, an architect and a worker each
* carry one.
*/
private static String callerTerminal(McpSyncServerExchange exchange) {
Object v = exchange.transportContext().get(CALLER_TERMINAL);
String s = v == null ? null : v.toString();
@@ -904,7 +908,8 @@ public final class FleetMcp {
* a BUSY interloper never claims a turn it did not win. {@code null} disables recording.
*/
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
Long timeoutMs, Runnable onAccepted, Set<String> profiles) {
Long timeoutMs, Runnable onAccepted, Set<String> profiles,
String callerTerminal) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -914,7 +919,7 @@ public final class FleetMcp {
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
try {
return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout);
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerTerminal), timeout);
} catch (HerdrException e) {
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
}
@@ -924,13 +929,15 @@ public final class FleetMcp {
* {@code fleet_send} carrying a {@code turnId}: the primary's answer to a worker's
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
* it resumes the same turn — surfaced to the primary identically to a normal send.
* {@code callerTerminal} must match the turn's recorded owner or this is refused.
*/
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs) {
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs,
String callerTerminal) {
if (isBlank(turnId) || isBlank(content)) {
return error("turnId and content are required to answer a worker's question");
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
return formatReply(messages.answer(turnId, content, timeout), timeout);
return formatReply(messages.answer(turnId, content, timeout, callerTerminal), timeout);
}
/**
@@ -979,6 +986,8 @@ public final class FleetMcp {
+ "\" and content set to your answer; the worker resumes the same turn.");
case STALE_TURN -> error("that question is no longer open — it timed out or was already "
+ "answered (turnId stale)");
case NOT_TURN_OWNER -> error("this turn belongs to a different delegation — only the caller "
+ "that opened it may answer it");
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
// fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call,
@@ -87,7 +87,7 @@ public final class MessageService {
/**
* The worker paused mid-turn to ask the primary a question (CB-205); {@code text} is the
* question and {@code turnId} correlates the answer. Not terminal — the primary answers with
* {@link #answer(String, String, long)} and the turn resumes.
* {@link #answer(String, String, long, String)} and the turn resumes.
*/
QUESTION,
/** Timed out after the message was delivered — the worker is still working. */
@@ -120,10 +120,16 @@ public final class MessageService {
/** Another send to this session was in flight for the whole window. */
BUSY,
/**
* An answer ({@link #answer(String, String, long)}) referenced a {@code turnId} that is no
* longer open — the worker's {@code fleet_ask} already timed out or was answered.
* An answer ({@link #answer(String, String, long, String)}) referenced a {@code turnId}
* that is no longer open — the worker's {@code fleet_ask} already timed out or was answered.
*/
STALE_TURN
STALE_TURN,
/**
* An answer ({@link #answer(String, String, long, String)}) named a {@code turnId} that is
* still open, but the answering caller is not the caller whose accepted delegation opened
* it. Distinct from {@link #STALE_TURN} so a refusal is never reported as a lapsed turn.
*/
NOT_TURN_OWNER
}
/**
@@ -133,7 +139,7 @@ public final class MessageService {
* {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION},
* else {@code null}
* @param turnId correlation id for a {@link Outcome#QUESTION} (answered via
* {@link #answer(String, String, long)}), else {@code null}
* {@link #answer(String, String, long, String)}), else {@code null}
*/
public record Reply(Outcome outcome, String text, String turnId) {
/** A reply with no correlation id (the common terminal outcomes). */
@@ -665,7 +671,7 @@ public final class MessageService {
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, TIMED_OUT_UNCONFIRMED, BUSY -> "timeout";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
case STALE_TURN, QUESTION -> null; // not a completed delegation
case STALE_TURN, QUESTION, NOT_TURN_OWNER -> null; // not a completed delegation
};
}
@@ -924,14 +930,17 @@ public final class MessageService {
/**
* Deliver {@code content} to {@code target} (a herdr {@code terminal_id}) and block until the
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses.
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. {@code callerTerminal}
* is the terminal of the caller making this call — {@code null} for the unnamed primary — and is
* recorded as the turn's owner, the only caller {@link #answer(String, String, long, String)} will
* later accept an answer from if the worker pauses mid-turn to ask.
*/
public Reply send(String target, String content, long timeoutMillis) {
return send(target, content, timeoutMillis, null);
public Reply send(String target, String content, long timeoutMillis, String callerTerminal) {
return send(target, content, timeoutMillis, null, callerTerminal);
}
/**
* As {@link #send(String, String, long)}, but with an accepted-delivery hook.
* As {@link #send(String, String, long, String)}, but with an accepted-delivery hook.
*
* <p>{@code onAccepted} is invoked exactly once, once this send has won {@code target}'s send
* lock and so become the <em>accepted target turn</em> — it runs <em>before</em> delivery is
@@ -942,12 +951,13 @@ public final class MessageService {
* acceptance means a concurrent sender that times out {@code BUSY} can never steal ownership it
* never earned. {@code null} disables the hook.
*/
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted) {
return send(target, content, timeoutMillis, onAccepted, null);
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerTerminal) {
return send(target, content, timeoutMillis, onAccepted, null, callerTerminal);
}
/** Run a send, optionally stopping an async task that teardown already failed before acceptance. */
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task) {
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task,
String callerTerminal) {
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
@@ -967,7 +977,7 @@ public final class MessageService {
// open race). Opening first also means a throwing onAccepted (fired before enqueue) or an
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
// failed send leaves no stale waiter behind.
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target, Rendezvous.Owner.of(callerTerminal));
// CB-640: this send now owns target's delivery, so any earlier stranded-reply or
// still-queued fact no longer describes the live state — clear both rather than let
// them outlive the send that supersedes them.
@@ -1146,19 +1156,30 @@ public final class MessageService {
* mid-turn (already picked up), so the answer flows back through its own open {@code fleet_ask}
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
* is unblocked so a reply that lands the instant it resumes is not lost.
*
* <p>{@code callerTerminal} is the terminal of the caller making this call — {@code null} for
* the unnamed primary. It is checked against the turn's recorded owner (the caller whose
* accepted delegation opened it, see {@link #send(String, String, long, String)} and
* {@link #sendAsync(String, String, Runnable, String)}) before anything else runs: a mismatch,
* including a turn with no owner on record at all, returns {@link Outcome#NOT_TURN_OWNER}
* without touching the rendezvous, the session lock, or any async task bookkeeping.
*/
public Reply answer(String turnId, String content, long timeoutMillis) {
public Reply answer(String turnId, String content, long timeoutMillis, String callerTerminal) {
String workerSession = rendezvous.askSession(turnId);
if (workerSession == null) {
return new Reply(Outcome.STALE_TURN, null); // the ask lapsed (timed out or already answered)
}
Rendezvous.Owner owner = rendezvous.askOwner(turnId);
if (!Rendezvous.Owner.permits(owner, callerTerminal)) {
return new Reply(Outcome.NOT_TURN_OWNER, null);
}
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(workerSession, _ -> new ReentrantLock());
if (!tryLock(lock, remainingMillis(deadlineNanos))) {
return new Reply(Outcome.BUSY, null);
}
try {
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession, owner);
// fleetd #575: this try used to open below, AFTER the Task lookup/registration and the
// STALE_TURN early return that follows it — so that return was covered only by a
// hand-rolled copy of the finally's own cleanup pair, not the finally itself. Widening the
@@ -1332,7 +1353,7 @@ public final class MessageService {
}
asyncExecutor.submit(() -> {
try {
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorTerminal);
if (result.outcome() == Outcome.QUESTION) {
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
// just after resolveQuestion wakes this thread.
@@ -60,7 +60,7 @@ public final class Rendezvous {
}
/** A worker's open mid-turn question: the worker session it belongs to and the answer future. */
private record AskWaiter(String session, CompletableFuture<String> answer) {
private record AskWaiter(String session, CompletableFuture<String> answer, Owner owner) {
}
/**
@@ -70,7 +70,33 @@ public final class Rendezvous {
public record AskTicket(String turnId, CompletableFuture<String> answer, boolean fresh) {
}
private final ConcurrentHashMap<String, CompletableFuture<Resolution>> waiters = new ConcurrentHashMap<>();
/**
* The caller whose accepted delegation opened a turn — the only caller allowed to answer it.
* A {@code null} terminal means the unnamed primary, an authenticated caller with no pane.
*/
public record Owner(String terminal) {
public static final Owner UNNAMED_PRIMARY = new Owner(null);
public static Owner of(String terminal) {
return terminal == null ? UNNAMED_PRIMARY : new Owner(terminal);
}
/**
* Whether {@code callerTerminal} matches {@code owner}. A {@code null} owner means no
* owner was ever recorded, and that state matches no caller, not even one whose own
* terminal is {@code null} — "no record" and "recorded as the unnamed primary" are
* different states.
*/
public static boolean permits(Owner owner, String callerTerminal) {
return owner != null && java.util.Objects.equals(owner.terminal(), callerTerminal);
}
}
/** A registered forward waiter together with the owner its delegation was opened under. */
private record ForwardWaiter(Owner owner, CompletableFuture<Resolution> future) {
}
private final ConcurrentHashMap<String, ForwardWaiter> waiters = new ConcurrentHashMap<>();
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */
private final ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
@@ -89,8 +115,17 @@ public final class Rendezvous {
* code a double open is impossible; this is a tripwire for the day that no longer holds.
*/
public CompletableFuture<Resolution> open(String session) {
return open(session, null);
}
/**
* Same as {@link #open(String)}, additionally recording {@code owner} as the caller whose
* delegation opened this waiter. A {@code null} owner records no owner at all — the
* fail-closed default {@link Owner#permits} refuses to everyone.
*/
public CompletableFuture<Resolution> open(String session, Owner owner) {
CompletableFuture<Resolution> waiter = new CompletableFuture<>();
CompletableFuture<Resolution> existing = waiters.putIfAbsent(session, waiter);
ForwardWaiter existing = waiters.putIfAbsent(session, new ForwardWaiter(owner, waiter));
if (existing != null) {
throw new IllegalStateException(
"rendezvous double-open for session " + session + " — a waiter is already registered");
@@ -105,7 +140,13 @@ public final class Rendezvous {
* successful {@code open} after a finished turn requires this close to have happened first).
*/
public void close(String session, CompletableFuture<Resolution> waiter) {
waiters.remove(session, waiter);
waiters.computeIfPresent(session, (s, w) -> w.future() == waiter ? null : w);
}
/** The owner recorded for {@code session}'s open waiter, or {@code null} if none is open. */
public Owner ownerOf(String session) {
ForwardWaiter w = waiters.get(session);
return w == null ? null : w.owner();
}
/** Whether a send is currently awaiting a resolution for {@code session}. */
@@ -119,7 +160,8 @@ public final class Rendezvous {
* send (see the CB-116 note above) rather than whichever send happens to be waiting when they fire.
*/
public CompletableFuture<Resolution> currentWaiter(String session) {
return waiters.get(session);
ForwardWaiter w = waiters.get(session);
return w == null ? null : w.future();
}
/**
@@ -147,7 +189,7 @@ public final class Rendezvous {
String turnId = openAsksBySession.computeIfAbsent(session, _ -> {
String newTurnId = session + "#" + askSeq.incrementAndGet();
CompletableFuture<String> answer = new CompletableFuture<>();
AskWaiter waiter = new AskWaiter(session, answer);
AskWaiter waiter = new AskWaiter(session, answer, ownerOf(session));
asks.put(newTurnId, waiter);
minted[0] = waiter;
return newTurnId;
@@ -183,6 +225,16 @@ public final class Rendezvous {
return w == null ? null : w.session();
}
/**
* The owner recorded for {@code turnId} when its ask turn was freshly opened — the caller
* whose delegation {@link #answerAsk} must match. {@code null} if {@code turnId} is unknown or
* lapsed, or if the ask opened with no forward waiter owner on record.
*/
public Owner askOwner(String turnId) {
AskWaiter w = asks.get(turnId);
return w == null ? null : w.owner();
}
/**
* Resolve a worker's blocked {@code fleet_ask} with the primary's {@code answer}, unblocking it
* to resume its turn.
@@ -245,7 +297,7 @@ public final class Rendezvous {
}
private boolean complete(String session, Resolution resolution) {
CompletableFuture<Resolution> waiter = waiters.get(session);
return waiter != null && waiter.complete(resolution);
ForwardWaiter waiter = waiters.get(session);
return waiter != null && waiter.future().complete(resolution);
}
}
@@ -663,6 +663,8 @@ public final class FleetApp {
if (!allow(ctx, routeAction("POST /sessions/{id}/message"), id)) {
return;
}
Principal caller = ctx.attribute(CALLER);
String callerTerminal = caller == null ? null : caller.terminal();
JsonNode body;
try {
body = mapper.readTree(ctx.body());
@@ -688,20 +690,19 @@ public final class FleetApp {
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
if (turnId != null && !turnId.isBlank()) {
writeReply(ctx, id, messages.answer(turnId, content, timeout), timeout);
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerTerminal), timeout);
return;
}
if (!wait) {
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
Principal caller = ctx.attribute(CALLER);
String ticket = messages.sendAsync(id, content, null, caller == null ? null : caller.terminal());
String ticket = messages.sendAsync(id, content, null, callerTerminal);
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
return;
}
try {
writeReply(ctx, id, messages.send(id, content, timeout), timeout);
writeReply(ctx, id, messages.send(id, content, timeout, callerTerminal), timeout);
} catch (HerdrException e) {
herdrError(ctx, e);
}
@@ -720,6 +721,10 @@ public final class FleetApp {
case STALE_TURN -> ctx.status(409).json(Map.of(
"sessionId", id, "error", "stale_turn",
"detail", "that question is no longer open (timed out or already answered)"));
case NOT_TURN_OWNER -> ctx.status(403).json(Map.of(
"sessionId", id, "error", "not_turn_owner",
"detail", "this turn belongs to a different delegation — only the caller that "
+ "opened it may answer it"));
case REPLIED, COMPLETED_UNREPLIED -> {
// replySource distinguishes a structured fleet_reply from the CB-106 completion
// fallback (a scrape of the worker's transcript when it finished without replying).
@@ -734,9 +739,10 @@ public final class FleetApp {
// silent fall-through. That is exactly the bug this ticket exists to fix:
// `default -> "done"` used to sit here and would have told a REST caller the
// delegation completed for TIMED_OUT_UNCONFIRMED, the one outcome where delivery
// is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION and STALE_TURN can never
// actually reach this inner switch — the outer switch above always dispatches
// them first — but they still need an arm to keep this switch exhaustive.
// is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN and
// NOT_TURN_OWNER can never actually reach this inner switch — the outer switch
// above always dispatches them first — but they still need an arm to keep this
// switch exhaustive.
"status", switch (reply.outcome()) {
case TIMED_OUT_WORKING -> "working";
case TIMED_OUT_QUEUED -> "queued";
@@ -746,7 +752,7 @@ public final class FleetApp {
case BUSY -> "busy";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN, NOT_TURN_OWNER -> "done"; // unreachable
},
"detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED
? "no reply within " + timeout + "ms; delivery is unconfirmed — the "
@@ -368,18 +368,18 @@ public final class SessionManager implements TurnListener {
/**
* Depth count plus the terminal/role a mid-teardown pane belongs to, for
* {@link #spawnedMemberRole}. The terminal/role come from whichever call into
* {@link #releaseWindow} first knew them — a losing {@link #releaseIfCurrent} CAS has no
* {@code known} session of its own, so it must not blank out what the winner already recorded.
* {@link #releaseWindow} first knew them: a call that finds the registry entry already gone
* passes a {@code null} session, and must not blank out what the first call recorded.
*/
private record Releasing(int depth, String terminalId, MemberRole role) {
private static Releasing enter(Releasing prior, MemberSession known) {
record Releasing(int depth, String terminalId, MemberRole role) {
static Releasing enter(Releasing prior, MemberSession known) {
int depth = (prior == null ? 0 : prior.depth()) + 1;
String terminalId = known != null ? known.terminalId() : prior == null ? null : prior.terminalId();
MemberRole role = known != null ? known.role() : prior == null ? null : prior.role();
return new Releasing(depth, terminalId, role);
}
private Releasing leave() {
Releasing leave() {
return depth <= 1 ? null : new Releasing(depth - 1, terminalId, role);
}
}
@@ -30,6 +30,22 @@ class AuthzTest {
assertTrue(Authz.isUnauthenticated(null));
}
/**
* {@code MessageService.answer}'s turn-ownership check treats a caller with no terminal as
* matching a turn recorded for the unnamed primary. That rule only stays safe because an
* unauthenticated caller — whose terminal is also {@code null} — never reaches {@code answer}
* at all: {@link #anonymousIsAuthorizedForNothing} already covers every action including
* {@code ANSWER}, but this test names the exact coupling so a future change to either side
* cannot drift without turning this test red.
*/
@Test
void anAnonymousCallerIsRefusedAnswerSoItCanNeverBeMistakenForTheUnnamedPrimary() {
assertFalse(Authz.permits(ANON, ANSWER, null),
"an anonymous caller, whose terminal is also null, must never reach answer() — the "
+ "turn-ownership check's null-terminal match for the unnamed primary owner "
+ "relies on this gate refusing it first");
}
@Test
void orchestrationBelongsToThePrimaryAlone() {
for (Authz.Action a : new Authz.Action[]{SPAWN, STOP, SEND, DRAIN}) {
@@ -77,7 +77,7 @@ class FleetMcpTest {
private void assertSendRoundTrips(String target, Set<String> profiles) throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles));
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -103,7 +103,7 @@ class FleetMcpTest {
void sendThenReplyRoundTrips() throws Exception {
// fleet_send blocks; fleet_reply resolves it with the worker's structured answer.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of()));
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null));
// Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now
// queues in the inbox if no waiter is open, which would break the round-trip).
@@ -181,7 +181,7 @@ class FleetMcpTest {
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
@@ -279,7 +279,7 @@ class FleetMcpTest {
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L));
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -288,6 +288,89 @@ class FleetMcpTest {
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
}
// --- fleetd #715: fleet_send{turnId} is gated on the caller that owns the turn -------------
/**
* A blocking {@code fleet_send} from one caller opens the turn; a {@code fleet_send{turnId}}
* from a different caller is refused as an error, and the real owner's answer still succeeds.
*/
@Test
void aDifferentCallersMcpAnswerIsRefusedForABlockingSendButTheRealOwnerSucceeds() throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner"));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T));
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> FleetMcp.ask(messages, T, "which config?", 5000L));
McpSchema.CallToolResult question = send.get(5, TimeUnit.SECONDS);
assertTrue(textOf(question).contains("[question]"), textOf(question));
String questionText = textOf(question);
String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
assertTrue(hijacked.isError(), "a caller that did not open this turn must get an error, not an answer");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
FleetMcp.reply(messages, T, Role.WORKER, "done");
assertEquals("done", textOf(answer.get(5, TimeUnit.SECONDS)));
}
/**
* Same hijack and control through the fire-and-poll ({@code sendAsync}) path: the owner comes
* from the ticket's recorded creator terminal, not from a caller threaded through a live call.
*/
@Test
void aDifferentCallersMcpAnswerIsRefusedForAnAsyncSendButTheRealOwnerSucceeds() throws Exception {
McpSchema.CallToolResult accepted =
FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), "term_owner");
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T));
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> FleetMcp.ask(messages, T, "which config?", 5000L));
MessageService.TaskView asking = messages.poll(ticket);
deadline = System.currentTimeMillis() + 3000;
while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
asking = messages.poll(ticket);
}
assertEquals(MessageService.Phase.ASKING, asking.phase());
String turnId = asking.turnId();
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
assertTrue(hijacked.isError(), "a caller that did not create this delegation must get an error");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
"a refused answer must not advance the async ticket's phase");
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
FleetMcp.reply(messages, T, Role.WORKER, "done");
assertEquals("done", textOf(answer.get(5, TimeUnit.SECONDS)));
}
@Test
void pollUnknownTicketIsAnError() {
McpSchema.CallToolResult res = FleetMcp.poll(messages, "task-999", null);
@@ -297,7 +380,7 @@ class FleetMcpTest {
@Test
void sendTimesOutWithAWorkingNote() {
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of());
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
}
@@ -313,7 +396,7 @@ class FleetMcpTest {
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
herdr.agentSendFailsWith("send_failed");
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of()));
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null));
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -333,13 +416,13 @@ class FleetMcpTest {
@Test
void sendRejectsMissingArgs() {
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError());
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of()).isError());
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null).isError());
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null).isError());
}
@Test
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"));
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null);
McpSchema.CallToolResult async = FleetMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
assertTrue(blocking.isError());
@@ -358,7 +441,7 @@ class FleetMcpTest {
assertSendRoundTrips("term_live_member", profiles);
// A herdr-owned pane outside the bridge roster cannot be classified at accept time.
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles);
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null);
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
}
@@ -450,7 +533,7 @@ class FleetMcpTest {
void askThenAnswerRoundTrips() throws Exception {
// The primary delegates and blocks; wait until its waiter is open before the worker asks.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of()));
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -472,7 +555,7 @@ class FleetMcpTest {
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
// The worker's ask returns the answer — it resumes the same turn.
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
@@ -498,7 +581,7 @@ class FleetMcpTest {
@Test
void answerToAStaleTurnIsAnError() {
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L);
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null);
assertTrue(res.isError());
assertTrue(textOf(res).contains("no longer open"), textOf(res));
}
@@ -1901,7 +1984,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));
() -> 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) {
@@ -1958,7 +2041,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));
() -> messages.answer(turnId, "config.yaml", 5000, "term_creator"));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -78,7 +78,7 @@ class MessageServiceTest {
}
private CompletableFuture<MessageService.Reply> sendAsync(String content, long timeoutMillis) {
return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis));
return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis, null));
}
private void awaitWaiting() throws InterruptedException {
@@ -188,7 +188,7 @@ class MessageServiceTest {
localRendezvous, new InMemoryReplyInbox());
String brief = "Implement the requested change. ".repeat(20);
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000));
CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000, (String) null));
long deadline = System.currentTimeMillis() + 2000;
while (!localRendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -341,7 +341,7 @@ class MessageServiceTest {
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null));
// The worker's ask returns the answer — it resumes the same turn.
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
@@ -377,7 +377,7 @@ class MessageServiceTest {
// The primary answers that one turnId; both asks unblock with the same answer.
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null));
MessageService.AskResult a1 = ask1.get(5, TimeUnit.SECONDS);
MessageService.AskResult a2 = ask2.get(5, TimeUnit.SECONDS);
@@ -472,7 +472,7 @@ class MessageServiceTest {
"the duplicate's own timeout elapses first");
CompletableFuture<MessageService.Reply> answered =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500));
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500, null));
MessageService.AskResult a = fresh.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome(),
@@ -507,7 +507,7 @@ class MessageServiceTest {
CompletableFuture<MessageService.Reply> lateAnswer = new CompletableFuture<>();
messages.setAskTimeoutRaceHookForTest(() ->
lateAnswer.complete(messages.answer(turnId, "too late", 500)));
lateAnswer.complete(messages.answer(turnId, "too late", 500, null)));
try {
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome());
@@ -539,7 +539,7 @@ class MessageServiceTest {
@Test
void answeringAnUnknownTurnIsStale() {
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
MessageService.Reply r = messages.answer(T + "#999", "too late", 500, null);
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
"an answer to a turn that never existed (or already lapsed) is stale, not a hang");
}
@@ -572,7 +572,7 @@ class MessageServiceTest {
messages.setAnswerAskLapseRaceHookForTest(() -> rendezvous.answerAsk(turnId, "raced in first"));
try {
MessageService.Reply r = messages.answer(turnId, "too late", 500);
MessageService.Reply r = messages.answer(turnId, "too late", 500, null);
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
"an ask already answered by the race must be seen as lapsed, not double-delivered");
} finally {
@@ -596,7 +596,7 @@ class MessageServiceTest {
void sendTimesOutBeforeDeliveryIsQueuedNotWorking() {
// Nothing ever delivers the message and nothing resolves the send, so the reply future
// times out with delivery still incomplete — the message is still queued for the worker.
MessageService.Reply r = messages.send(T, "never delivered", 50);
MessageService.Reply r = messages.send(T, "never delivered", 50, null);
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome(),
"an undelivered send that times out is still queued, not working");
assertNull(r.text());
@@ -605,7 +605,7 @@ class MessageServiceTest {
@Test
void sendTimesOutAfterDeliveryIsStillWorking() throws Exception {
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 300));
CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 300, null));
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver — the delivered future now completes
injector.onStatus(T, AgentStatus.WORKING); // worker starts but never replies
@@ -630,7 +630,7 @@ class MessageServiceTest {
void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
try {
MessageService.Reply reply = messages.send(T, "race delivery", 50);
MessageService.Reply reply = messages.send(T, "race delivery", 50, null);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
"cancel reporting DELIVERED means the worker received the timed-out message");
@@ -652,7 +652,7 @@ class MessageServiceTest {
void sendTimesOutWithAttemptedDeliveryReportsUnconfirmedNotQueued() throws Exception {
herdr.agentSendFailsWith("send_failed");
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150));
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150, null));
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt → ATTEMPTED
@@ -679,7 +679,7 @@ class MessageServiceTest {
// The primary answers, unblocking the worker; but the worker never sends the follow-up
// fleet_reply, so the answering send rides out its short window as still-working.
MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200);
MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200, null);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answer.outcome(),
"an answered worker that never replies times out as still working");
@@ -715,7 +715,7 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null));
ask.get(5, TimeUnit.SECONDS); // worker resumed with the answer
awaitWaiting(); // the answering call has (re)opened its own forward waiter
@@ -723,7 +723,7 @@ class MessageServiceTest {
MessageService.Reply done = answer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.REPLIED, done.outcome());
MessageService.Reply probe = messages.send(T, "probe after normal reply", 300);
MessageService.Reply probe = messages.send(T, "probe after normal reply", 300, null);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released after a normal REPLIED answer(), or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
@@ -756,13 +756,13 @@ class MessageServiceTest {
// The primary answers, unblocking the worker; the worker never sends its follow-up
// fleet_reply, so the answering call rides out its short window as still-working.
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200, null));
MessageService.Reply answered = answer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answered.outcome(),
"an answered worker that never replies times out as still working");
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300);
MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300, null);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released after a TIMED_OUT_WORKING answer(), or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
@@ -791,7 +791,7 @@ class MessageServiceTest {
new java.util.concurrent.atomic.AtomicReference<>();
Thread answerer = new Thread(() -> {
try {
messages.answer(turnId, "config.yaml", 5000);
messages.answer(turnId, "config.yaml", 5000, null);
caught.set(new AssertionError("expected answer() to throw"));
} catch (Throwable t) {
caught.set(t);
@@ -813,7 +813,7 @@ class MessageServiceTest {
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300);
MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300, null);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released when answer()'s reply future fails exceptionally, "
+ "or this bounded follow-up send would come back BUSY instead of timing out on its own work");
@@ -840,7 +840,7 @@ class MessageServiceTest {
new java.util.concurrent.atomic.AtomicReference<>();
Thread answerer = new Thread(() -> {
try {
messages.answer(turnId, "config.yaml", 5000);
messages.answer(turnId, "config.yaml", 5000, null);
caught.set(new AssertionError("expected answer() to throw"));
} catch (Throwable t) {
caught.set(t);
@@ -859,7 +859,7 @@ class MessageServiceTest {
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after interruption", 300);
MessageService.Reply probe = messages.send(T, "probe after interruption", 300, null);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released when answer()'s wait is interrupted, or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
@@ -967,11 +967,11 @@ class MessageServiceTest {
@Test
void concurrentSendToSameSessionWhileFirstHoldsItIsBusy() throws Exception {
CompletableFuture<MessageService.Reply> first =
CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000));
CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000, null));
awaitWaiting(); // the first send now holds the session lock, blocked on its reply
// A second send to the SAME session cannot take the lock within its short window.
MessageService.Reply busy = messages.send(T, "second", 100);
MessageService.Reply busy = messages.send(T, "second", 100, null);
assertEquals(MessageService.Outcome.BUSY, busy.outcome(),
"a second send while another holds the session is busy, not a hang");
assertNull(busy.text());
@@ -1001,13 +1001,13 @@ class MessageServiceTest {
PrimaryRegistry reg = new PrimaryRegistry(null);
// L accepts a delegation to W: the send wins the lock and queues delivery → L is recorded.
CompletableFuture<MessageService.Reply> first = CompletableFuture.supplyAsync(
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L)));
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L), null));
awaitWaiting();
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
"an accepted send owns the delegation");
// A attempts W while L holds it → BUSY (lock never taken) → its hook never fires.
MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A));
MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A), null);
assertEquals(MessageService.Outcome.BUSY, busy.outcome());
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
"a BUSY send must not steal the delegator ownership it never earned");
@@ -1028,7 +1028,7 @@ class MessageServiceTest {
void anAcceptedSendAfterThePriorOwnerFinishesBecomesTheNewOwner() throws Exception {
PrimaryRegistry reg = new PrimaryRegistry(null);
CompletableFuture<MessageService.Reply> first = CompletableFuture.supplyAsync(
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L)));
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L), null));
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
@@ -1039,7 +1039,7 @@ class MessageServiceTest {
// L finished; A's later accepted send takes the delegation over.
CompletableFuture<MessageService.Reply> second = CompletableFuture.supplyAsync(
() -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A)));
() -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A), null));
awaitWaiting();
assertEquals(LEAD_A, reg.nudgeTargetFor(T).orElseThrow(),
"an accepted send after the owner finished becomes the new delegator");
@@ -1060,7 +1060,7 @@ class MessageServiceTest {
void answeringAnAskDoesNotRewriteDelegatorOwnership() throws Exception {
PrimaryRegistry reg = new PrimaryRegistry(null);
CompletableFuture<MessageService.Reply> send = CompletableFuture.supplyAsync(
() -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L)));
() -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L), null));
awaitWaiting();
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "L owns the delegation");
@@ -1075,7 +1075,7 @@ class MessageServiceTest {
// L answers the ask on the same turn; the answer path must not touch ownership.
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000, null));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // the answering send reopened its forward waiter
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
@@ -1095,7 +1095,7 @@ class MessageServiceTest {
void aThrowingAcceptedHookLeavesNoStaleWaiterOrQueuedOrphan() {
assertThrows(IllegalStateException.class,
() -> messages.send(T, "doomed", 500,
() -> { throw new IllegalStateException("ownership hook failed"); }),
() -> { throw new IllegalStateException("ownership hook failed"); }, null),
"a throwing ownership hook fails the send loudly");
assertFalse(rendezvous.isWaiting(T), "the failed send must not leave a stale rendezvous waiter");
@@ -1223,7 +1223,7 @@ class MessageServiceTest {
@Test
void abandonFailsASendThatIsStillWaitingOnAReleasedSession() throws Exception {
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000, null));
awaitWaiting();
assertTrue(messages.abandon(T, "session released"), "a live waiter is abandoned");
@@ -1243,7 +1243,7 @@ class MessageServiceTest {
@Test
void abandonDoesNotOverwriteAnAlreadyResolvedSend() throws Exception {
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000, null));
awaitWaiting();
assertTrue(rendezvous.resolve(T, "the real answer"));
@@ -1374,7 +1374,7 @@ class MessageServiceTest {
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
() -> messages.answer(asking.turnId(), "config.yaml", 5000, null));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1414,7 +1414,7 @@ class MessageServiceTest {
// same turnId must see it as lapsed rather than resolving a question nobody is waiting on.
assertEquals(MessageService.AskOutcome.TIMED_OUT, ask.get(5, TimeUnit.SECONDS).outcome());
assertEquals(MessageService.Outcome.STALE_TURN,
messages.answer(asking.turnId(), "config.yaml", 200).outcome());
messages.answer(asking.turnId(), "config.yaml", 200, null).outcome());
}
@Test
@@ -1451,7 +1451,7 @@ class MessageServiceTest {
// The primary answers, but its own bounded wait for the worker's resumed turn is short and
// expires before the worker (still genuinely working) gets back to it.
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150, null);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait gives up before the worker finishes resuming");
@@ -1480,7 +1480,7 @@ class MessageServiceTest {
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150, null);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
@@ -1534,7 +1534,7 @@ class MessageServiceTest {
messages.setFinishAsyncTaskRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000, null));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
@@ -1640,7 +1640,7 @@ class MessageServiceTest {
String turnId = asking.turnId();
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000, null));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer(),
"the worker's own ask() call must have already unblocked with the primary's answer "
+ "before we force the race below");
@@ -1687,7 +1687,7 @@ class MessageServiceTest {
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150);
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150, null);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait must give up first, leaving turnId stamped with no "
@@ -1753,7 +1753,7 @@ class MessageServiceTest {
MessageService.TaskView asking = messages.poll(first);
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
() -> messages.answer(asking.turnId(), "config.yaml", 5000, null));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1787,7 +1787,7 @@ class MessageServiceTest {
// The primary answers it — answer() resumes the turn and blocks for what comes next.
CompletableFuture<MessageService.Reply> answer1 = CompletableFuture.supplyAsync(
() -> messages.answer(asking1.turnId(), "a1", 5000));
() -> messages.answer(asking1.turnId(), "a1", 5000, null));
assertEquals("a1", ask1.get(5, TimeUnit.SECONDS).answer());
// Still in the SAME resumed turn — before replying — the worker asks again.
@@ -1808,7 +1808,7 @@ class MessageServiceTest {
// The primary answers the second question; the worker finally sends its real fleet_reply.
CompletableFuture<MessageService.Reply> answer2 = CompletableFuture.supplyAsync(
() -> messages.answer(turnId2, "a2", 5000));
() -> messages.answer(turnId2, "a2", 5000, null));
assertEquals("a2", ask2.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1908,7 +1908,7 @@ class MessageServiceTest {
assertEquals(asking.turnId(), pending.turnId());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
() -> messages.answer(asking.turnId(), "config.yaml", 5000, null));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertNull(messages.pendingAsk(T, null), "an answered question is no longer pending");
awaitWaiting();
@@ -1943,7 +1943,7 @@ class MessageServiceTest {
assertEquals("which config file?", unnamed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
() -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_creator"));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -1971,13 +1971,81 @@ class MessageServiceTest {
"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));
() -> 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
@@ -2234,7 +2302,7 @@ class MessageServiceTest {
// Clean up the still-open ask so the background thread does not linger past the test.
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000, null));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
@@ -2260,7 +2328,7 @@ class MessageServiceTest {
.filter(c -> c.method().equals("agent.prompt")).count();
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000, null));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
@@ -2653,7 +2721,7 @@ class MessageServiceTest {
void hasQueuedDeliveryIsTrueAfterAnUndeliveredSendTimesOut() {
// Nothing ever delivers the message (never goes IDLE/BLOCKED), so the send times out with
// TIMED_OUT_QUEUED — same setup as sendTimesOutBeforeDeliveryIsQueuedNotWorking above.
MessageService.Reply r = messages.send(T, "never delivered", 50);
MessageService.Reply r = messages.send(T, "never delivered", 50, null);
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
assertTrue(messages.hasQueuedDelivery(T),
@@ -2665,7 +2733,7 @@ class MessageServiceTest {
@Test
void hasQueuedDeliveryClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50, null).outcome());
assertTrue(messages.hasQueuedDelivery(T));
// A fresh send accepts delivery (opens its own waiter) — the stale queued fact is cleared.
@@ -2682,7 +2750,7 @@ class MessageServiceTest {
@Test
void hasQueuedDeliveryClearsOnAbandon() {
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50, null).outcome());
assertTrue(messages.hasQueuedDelivery(T));
messages.abandon(T, "session released");
@@ -0,0 +1,187 @@
package dev.ltms.fleet.msg;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.Stream;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Pins that no file under {@code src/main/java} calls the fail-open
* {@link Rendezvous#open(String)} overload. That overload opens a forward waiter with no
* recorded {@link Rendezvous.Owner}, so the turn it opens can never be answered by anyone,
* not even the unnamed primary. Every production caller must go through
* {@link Rendezvous#open(String, Rendezvous.Owner)} and record an explicit owner.
*
* <p>This reads each file's own source text rather than reflecting on compiled bytecode, because
* the risk is a future one-word edit at a call site, not a missing overload.
*
* <p>The scan below finds a violation by its receiver, {@code rendezvous.open(}, classified by
* argument count. The one test method here points that exact scanner at a file known to hold
* many real one-argument calls before it ever looks at production, so a scanner that stops
* matching fails loudly on the known-positive case instead of leaving a clean production result
* looking like evidence it never produced.
*/
class RendezvousOpenUsageTest {
private static final Path PRODUCTION_SOURCE = Path.of("src/main/java");
private static final Path KNOWN_TEST_CALLER =
Path.of("src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java");
/**
* A pattern that cannot find a known one-argument {@code rendezvous.open(} call would also
* find none in production -- not because production is clean, but because the pattern does
* not match the text it is supposed to catch. That failure mode is exactly what let a
* {@code \b}-based regex read "no callers" under {@code git grep -E} when 81 real ones
* existed: {@code git grep} does not treat {@code \b} as a word boundary, so the pattern
* silently matched nothing anywhere, clean code and real calls alike. This test runs the
* known-positive check first, with the same scanning method the production check then
* depends on, so that mistake fails loudly here instead of reading as a clean result.
*/
@Test
void noProductionFileCallsTheSingleArgumentOpenOverload() throws IOException {
List<String> knownCalls = new ArrayList<>();
int knownFilesScanned = scanForOneArgOpenCalls(KNOWN_TEST_CALLER, knownCalls);
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "found
// nothing" having looked at nothing.
assertTrue(knownFilesScanned > 0, "control failed: the scan under " + KNOWN_TEST_CALLER
+ " visited zero .java files -- the path is wrong, so neither result below proves "
+ "anything");
// CONTROL: the scanner actually finds real one-argument rendezvous.open( calls when
// pointed at a file known to hold many. If this is not satisfied, the matching logic
// itself is broken, and the production result below is the scanner failing silently,
// not production code actually being clean.
assertTrue(knownCalls.size() >= 60, "control failed: the scanner found only "
+ knownCalls.size() + " one-argument rendezvous.open( call(s) in " + KNOWN_TEST_CALLER
+ ", which is known to hold many -- the matching logic itself is broken: " + knownCalls);
List<String> violations = new ArrayList<>();
int filesScanned = scanForOneArgOpenCalls(PRODUCTION_SOURCE, violations);
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "no
// violations found" having looked at nothing.
assertTrue(filesScanned > 0, "control failed: the scan under " + PRODUCTION_SOURCE
+ " visited zero .java files -- the path is wrong, so the absence of violations "
+ "below proves nothing");
assertTrue(violations.isEmpty(), "found a call to the fail-open Rendezvous.open(String) "
+ "overload, which opens a forward waiter with no recorded owner -- record an "
+ "explicit Rendezvous.Owner through open(String, Owner) instead: " + violations);
}
private static int scanForOneArgOpenCalls(Path root, List<String> sites) throws IOException {
List<Path> files = javaFiles(root);
for (Path file : files) {
scanFileForOpenCalls(file, sites);
}
return files.size();
}
private static List<Path> javaFiles(Path root) throws IOException {
try (Stream<Path> paths = Files.walk(root)) {
return paths.filter(p -> p.toString().endsWith(".java")).toList();
}
}
private static void scanFileForOpenCalls(Path file, List<String> sites) throws IOException {
String source = Files.readString(file);
String needle = "rendezvous.open(";
int from = 0;
int idx;
while ((idx = source.indexOf(needle, from)) >= 0) {
int argsStart = idx + needle.length();
String args = extractBalancedArgs(source, argsStart, file, idx);
int closeParenIndex = argsStart + args.length();
from = closeParenIndex + 1;
if (args.isBlank()) {
continue; // Rendezvous has no zero-argument open() -- a javadoc "rendezvous.open()" mention, not a call
}
if (topLevelCommaCount(args) == 0) {
sites.add(file + ":" + lineOf(source, idx) + " -- rendezvous.open(" + args.trim() + ")");
}
}
}
/**
* The text between {@code rendezvous.open(} and its matching close paren: balanced over
* nested calls, and never split by a paren or comma sitting inside a string or char literal.
*/
private static String extractBalancedArgs(String source, int start, Path file, int callIndex) {
int depth = 1;
boolean inString = false;
boolean inChar = false;
int i = start;
while (i < source.length()) {
char c = source.charAt(i);
if (inString) {
if (c == '\\') { i += 2; continue; }
if (c == '"') inString = false;
} else if (inChar) {
if (c == '\\') { i += 2; continue; }
if (c == '\'') inChar = false;
} else if (c == '"') {
inString = true;
} else if (c == '\'') {
inChar = true;
} else if (c == '(') {
depth++;
} else if (c == ')') {
depth--;
if (depth == 0) return source.substring(start, i);
}
i++;
}
throw new IllegalStateException(
"unbalanced parentheses scanning " + file + ":" + lineOf(source, callIndex));
}
/**
* Commas at paren/bracket/brace depth zero, skipping string and char literals -- the argument
* separators a human reader would see, not every comma character in the text.
*/
private static int topLevelCommaCount(String args) {
int depth = 0;
int commas = 0;
boolean inString = false;
boolean inChar = false;
int i = 0;
while (i < args.length()) {
char c = args.charAt(i);
if (inString) {
if (c == '\\') { i += 2; continue; }
if (c == '"') inString = false;
} else if (inChar) {
if (c == '\\') { i += 2; continue; }
if (c == '\'') inChar = false;
} else if (c == '"') {
inString = true;
} else if (c == '\'') {
inChar = true;
} else if (c == '(' || c == '[' || c == '{') {
depth++;
} else if (c == ')' || c == ']' || c == '}') {
depth--;
} else if (c == ',' && depth == 0) {
commas++;
}
i++;
}
return commas;
}
private static int lineOf(String source, int index) {
int line = 1;
for (int i = 0; i < index; i++) {
if (source.charAt(i) == '\n') line++;
}
return line;
}
}
@@ -132,4 +132,80 @@ class RendezvousTest {
assertNull(rendezvous.askSession(t.turnId()), "a closed ask is forgotten");
assertFalse(rendezvous.answerAsk(t.turnId(), "late"), "a closed ask can no longer be answered");
}
// ── fleetd #715: turn ownership ────────────────────────────────────────────────────────────
@Test
void openWithNoOwnerRecordsNoOwnerAndOpenWithAnOwnerRecordsIt() {
rendezvous.open(W);
assertNull(rendezvous.ownerOf(W), "the no-owner overload records no owner at all");
rendezvous.close(W, rendezvous.currentWaiter(W));
rendezvous.open(W, Rendezvous.Owner.of("term_lead"));
assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.ownerOf(W));
}
@Test
void ownerOfIsNullWhenNoWaiterIsOpen() {
assertNull(rendezvous.ownerOf(W), "no waiter open means no owner to report");
}
@Test
void openAskStampsTheForwardWaitersOwnerOntoTheFreshTurnOnly() {
rendezvous.open(W, Rendezvous.Owner.of("term_lead"));
Rendezvous.AskTicket fresh = rendezvous.openAsk(W);
assertTrue(fresh.fresh());
assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.askOwner(fresh.turnId()),
"a freshly-opened ask copies the forward waiter's current owner");
Rendezvous.AskTicket coalesced = rendezvous.openAsk(W);
assertFalse(coalesced.fresh());
assertEquals(fresh.turnId(), coalesced.turnId());
assertEquals(Rendezvous.Owner.of("term_lead"), rendezvous.askOwner(coalesced.turnId()),
"a coalesced duplicate ask rides the fresh owner's turn, unchanged");
}
@Test
void aSecondAskAfterTheFirstClosesStampsWhateverOwnerIsOpenAtThatLaterMoment() {
rendezvous.open(W, Rendezvous.Owner.of("term_lead"));
Rendezvous.AskTicket first = rendezvous.openAsk(W);
rendezvous.closeAsk(first.turnId());
// The forward waiter is reopened under a different owner before the second ask — mirrors
// answer() reopening with the owner it already checked, which can differ turn to turn.
rendezvous.close(W, rendezvous.currentWaiter(W));
rendezvous.open(W, Rendezvous.Owner.of("term_other"));
Rendezvous.AskTicket second = rendezvous.openAsk(W);
assertTrue(second.fresh());
assertEquals(Rendezvous.Owner.of("term_other"), rendezvous.askOwner(second.turnId()),
"a freshly-opened ask always copies whatever owner is open right now, not a stale one");
}
@Test
void askOwnerIsNullForAnUnknownOrLapsedTurn() {
assertNull(rendezvous.askOwner("no-such#1"));
Rendezvous.AskTicket t = rendezvous.openAsk(W);
rendezvous.closeAsk(t.turnId());
assertNull(rendezvous.askOwner(t.turnId()), "a closed ask no longer reports an owner");
}
@Test
void ownerPermitsIsFailClosedOnARecordAndThreeStatesAreDistinct() {
assertFalse(Rendezvous.Owner.permits(null, null),
"no owner on record refuses even a caller with no terminal");
assertFalse(Rendezvous.Owner.permits(null, "term_a"),
"no owner on record refuses a terminal-bearing caller too");
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, null),
"the unnamed primary owner matches a caller with no terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, "term_a"),
"the unnamed primary owner does not match a terminal-bearing caller");
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_a"),
"a named owner matches the same terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_b"),
"a named owner refuses a different terminal");
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), null),
"a named owner refuses the unnamed primary");
}
}
@@ -343,7 +343,7 @@ class FleetAppAuthTest {
// 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));
() -> 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) {
@@ -358,6 +358,154 @@ class FleetAppAuthTest {
}
}
/**
* {@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
@@ -2558,78 +2558,38 @@ class SessionManagerTest {
}
/**
* {@code Releasing} is private, so its {@code enter}/{@code leave} are reached through
* reflection. Each is called directly, never through {@link SessionManager#release} or
* {@link SessionManager#spawnedMemberRole}, so a test naming one of them exercises only that
* one — a regression in the other can never hide behind it.
* {@code Releasing.enter} and {@code Releasing.leave} are called directly here, never through
* {@link SessionManager#release} or {@link SessionManager#spawnedMemberRole}, so a test naming
* one of them exercises only that one — a regression in the other can never hide behind it.
*/
private static Object releasingOf(int depth, String terminalId, MemberRole role) throws ReflectiveOperationException {
Class<?> cls = Class.forName("dev.ltms.fleet.session.SessionManager$Releasing");
java.lang.reflect.Constructor<?> ctor = cls.getDeclaredConstructor(int.class, String.class, MemberRole.class);
ctor.setAccessible(true);
return ctor.newInstance(depth, terminalId, role);
}
private static Object releasingEnter(Object prior, MemberSession known) throws ReflectiveOperationException {
Class<?> cls = Class.forName("dev.ltms.fleet.session.SessionManager$Releasing");
java.lang.reflect.Method m = cls.getDeclaredMethod("enter", cls, MemberSession.class);
m.setAccessible(true);
return m.invoke(null, prior, known);
}
private static Object releasingLeave(Object releasing) throws ReflectiveOperationException {
java.lang.reflect.Method m = releasing.getClass().getDeclaredMethod("leave");
m.setAccessible(true);
return m.invoke(releasing);
}
private static int depthOf(Object releasing) throws ReflectiveOperationException {
java.lang.reflect.Method m = releasing.getClass().getDeclaredMethod("depth");
m.setAccessible(true);
return (int) m.invoke(releasing);
}
private static String terminalIdOf(Object releasing) throws ReflectiveOperationException {
java.lang.reflect.Method m = releasing.getClass().getDeclaredMethod("terminalId");
m.setAccessible(true);
return (String) m.invoke(releasing);
}
private static MemberRole roleOf(Object releasing) throws ReflectiveOperationException {
java.lang.reflect.Method m = releasing.getClass().getDeclaredMethod("role");
m.setAccessible(true);
return (MemberRole) m.invoke(releasing);
}
@Test
void releasingLeaveStepsDownADepthGreaterThanOneInsteadOfRemovingIt() throws ReflectiveOperationException {
Object depthTwo = releasingOf(2, "term_a", MemberRole.DEV);
void releasingLeaveStepsDownADepthGreaterThanOneInsteadOfRemovingIt() {
SessionManager.Releasing depthTwo = new SessionManager.Releasing(2, "term_a", MemberRole.DEV);
Object afterLeave = releasingLeave(depthTwo);
SessionManager.Releasing afterLeave = depthTwo.leave();
assertNotNull(afterLeave,
"depth 2 means another release of the SAME pane is still mid-teardown; leave() must "
+ "step the depth down, never remove the marker outright — removing it here "
+ "is what a plain Set would do, and would reopen the window the still-in-"
+ "flight release is relying on staying closed");
assertEquals(1, depthOf(afterLeave));
assertEquals("term_a", terminalIdOf(afterLeave));
assertEquals(MemberRole.DEV, roleOf(afterLeave));
assertEquals(1, afterLeave.depth());
assertEquals("term_a", afterLeave.terminalId());
assertEquals(MemberRole.DEV, afterLeave.role());
}
@Test
void releasingEnterPreservesThePriorTerminalWhenTheOverlappingCallHasNoSessionOfItsOwn()
throws ReflectiveOperationException {
Object prior = releasingOf(1, "term_a", MemberRole.DEV);
void releasingEnterPreservesThePriorTerminalWhenTheOverlappingCallHasNoSessionOfItsOwn() {
SessionManager.Releasing prior = new SessionManager.Releasing(1, "term_a", MemberRole.DEV);
Object afterEnter = releasingEnter(prior, null);
SessionManager.Releasing afterEnter = SessionManager.Releasing.enter(prior, null);
assertEquals(2, depthOf(afterEnter), "depth still increments whether or not this enter knows its session");
assertEquals("term_a", terminalIdOf(afterEnter),
"a losing CAS (the releaseIfCurrent race fleetd #702 is about) has no session of "
+ "its own to pass as known, and must not blank out the terminal the winning "
assertEquals(2, afterEnter.depth(), "depth still increments whether or not this enter knows its session");
assertEquals("term_a", afterEnter.terminalId(),
"an overlapping release that finds the registry entry already gone has no session of "
+ "its own to pass as known, and must not blank out the terminal the first "
+ "call already recorded — that terminal is what spawnedMemberRole matches "
+ "against");
assertEquals(MemberRole.DEV, roleOf(afterEnter));
assertEquals(MemberRole.DEV, afterEnter.role());
}
}