Compare commits

...

15 Commits

Author SHA1 Message Date
Dai Ha 1fc9e85bf1 fleetd #729: fold a per-boot nonce into every turnId
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 1m57s
askSeq restarts at 0 on every daemon boot, so a turnId (session#n)
minted by one Rendezvous instance could be minted again by a later
instance and resolve to an unrelated ask. Fold a per-instance nonce
into the mint, the same way #719 fixed MessageService's ticket ids.
2026-10-04 18:53:05 +02:00
Dai Ha 809b7d9b20 Merge remote-tracking branch 'origin/worker/727-ee14ed-3'
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 47s
CI / build (push) Failing after 1m56s
2026-10-04 18:29:28 +02:00
Dai Ha 8cf7215d56 Merge remote-tracking branch 'origin/worker/719-bdd95e-4'
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 1m58s
2026-10-04 18:19:55 +02:00
Dai Ha 1ef93e57cc fleetd #727: give a lead launch the three protections every member spawn gets
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m59s
Extract checkPaneCommandFits, the agent_pane_busy retry, and the agent_name_taken
retry out of HerdrPeerLauncher into a shared dev.ltms.fleet.herdr.ResilientAgentLaunch,
and route LeadLauncher.launch through the same seam instead of a bare agents.start
call. The lead's agent name now carries a per-process nonce and a per-start sequence
number (like a member's), so a stale agent_name_taken from an earlier crashed session
no longer blocks a legitimate relaunch outright.
2026-10-04 18:14:06 +02:00
Dai Ha cf0c9b9316 fleetd #719: make the foreign-id test reach a colliding sequence number
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 1m49s
other.sendAsync had never been called, so other's tasks map was empty and
poll(ticket) returned null regardless of whether the nonce existed — the
test passed against an empty map, not against a colliding id. Mint once on
other so it reaches the same sequence number as the first instance, making
the test exercise the actual collision the nonce guards against.
2026-10-04 18:14:01 +02:00
Dai Ha 337dbd491e fleetd #719: fold a per-boot nonce into every ticket id
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 2m0s
ticketSeq restarted at zero on every daemon boot with no persistence, so a
ticket id minted in one boot could be reused by a later boot and resolve to
an unrelated Task instead of failing to resolve at all. Mint each
MessageService instance's own short nonce once and fold it into every ticket
(task-<nonce>-<n>), so an id from one instance can never match another's id
space.

Adds a disjoint-id-space test and a foreign-instance-ticket test (with the
positive control) in MessageServiceTest.
2026-10-04 18:05:06 +02:00
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
17 changed files with 1253 additions and 267 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
@@ -0,0 +1,146 @@
package dev.ltms.fleet.herdr;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.function.IntFunction;
/**
* The three protections every {@code agent.start} caller needs against herdr's pane-typed launch
* surface (fleetd #220, #727): a byte-limit check on the assembled command line, a bounded retry
* on {@code agent_pane_busy} (the target pane's shell has not reached its prompt yet), and a
* bounded retry on {@code agent_name_taken} with a fresh name each attempt. One implementation —
* every caller of {@code agent.start}, lead or member, goes through this seam rather than carrying
* its own copy.
*/
public final class ResilientAgentLaunch {
private static final Logger log = LoggerFactory.getLogger(ResilientAgentLaunch.class);
private ResilientAgentLaunch() {
}
/**
* The pty line buffer herdr types a launch command into: BSD/macOS {@code MAX_CANON}. Not a
* fleetd choice and not configurable — see {@link #checkFits}.
*/
public static final int PANE_COMMAND_BYTE_LIMIT = 1024;
/** Per-argument allowance for the separating space and a shell quote pair fleetd cannot see. */
private static final int QUOTING_OVERHEAD_PER_ARG = 3;
/** herdr rejects a duplicate agent {@code name}; a caller retries a bumped name this many times. */
public static final int NAME_RETRIES = 8;
/**
* Retries for {@code agent.start} against a seed pane whose shell has not reached its prompt
* yet — {@code tab.create}/{@code pane.split} return as soon as the pane exists, and herdr
* refuses to start an agent in a pane that is not "an available shell" ({@code agent_pane_busy}).
*/
public static final int SHELL_READY_RETRIES = 20;
/** Raised by {@link #checkFits} when the assembled command cannot fit the pane line. */
public static final class TooLargeException extends RuntimeException {
public TooLargeException(String message) {
super(message);
}
}
/**
* Verify the assembled launch line fits the pane herdr types it into. herdr does not exec the
* launch command — it TYPES it into the pane as one line, and a pty line buffer holds only
* {@value #PANE_COMMAND_BYTE_LIMIT} bytes. Everything past that byte is dropped with no error
* anywhere: herdr answers "agent started", the backend exits on the mangled argument it was
* handed, the pane closes, and the only symptom is a readiness timeout with no reason. That is
* how fleetd #214 broke every claude-code spawn — one 50-byte flag pushed a 978-byte command to
* 1028, and the tail that got cut was {@code --autocompact 250000}.
*
* <p>So measure it here and refuse, loudly and immediately, rather than start something that
* cannot work. The estimate is deliberately conservative: fleetd cannot see herdr's quoting, so
* every argument is charged its own bytes plus a separator and a quote pair. An over-estimate
* costs a clear error at a length that was already unsafe; an under-estimate would let the
* silent truncation back in.
*
* @param label names the launch in the refusal message (a profile name)
* @param argv the full argv, including the executable at index 0
* @throws TooLargeException naming the limit, the estimate, and the longest argument
*/
public static void checkFits(String label, List<String> argv) {
int bytes = 0;
String longest = null;
int longestBytes = 0;
for (String arg : argv) {
int argBytes = arg == null ? 0 : arg.getBytes(StandardCharsets.UTF_8).length;
bytes += argBytes + QUOTING_OVERHEAD_PER_ARG;
if (argBytes > longestBytes) {
longestBytes = argBytes;
longest = arg;
}
}
if (bytes <= PANE_COMMAND_BYTE_LIMIT) {
return;
}
String culprit = longest == null ? "<none>"
: longest.substring(0, Math.min(longest.length(), 60)) + (longest.length() > 60 ? "…" : "");
throw new TooLargeException(
"launch command for " + label + " is about " + bytes + " bytes, over the "
+ PANE_COMMAND_BYTE_LIMIT + "-byte limit of the pane line herdr types it into. "
+ "The pty would drop the tail silently and the backend would exit on a mangled "
+ "argument. Longest argument is " + longestBytes + " bytes: " + culprit
+ " — move it off the command line (a file flag) or shorten it.");
}
/**
* Start an agent into {@code paneId}, retrying {@code agent_pane_busy} up to {@code retries}
* times with {@code sleeper} run between attempts.
*
* @throws HerdrException the last {@code agent_pane_busy} failure once {@code retries} is
* spent, or immediately for any other herdr failure
*/
public static Agent startAwaitingShellPrompt(AgentControl agents, String name, String kind,
List<String> args, String paneId,
int retries, Runnable sleeper) {
HerdrException busy = null;
for (int attempt = 0; attempt < retries; attempt++) {
try {
return agents.start(name, kind, args, paneId);
} catch (HerdrException e) {
if (!"agent_pane_busy".equals(e.code())) throw e;
log.debug("pane {} not at its shell prompt yet, retrying agent.start", paneId);
busy = e;
sleeper.run();
}
}
throw busy;
}
/**
* Start an agent under a freshly generated name each attempt, retrying {@code agent_name_taken}
* up to {@code nameRetries} times — herdr refuses a duplicate {@code name} outright, so a stale
* registry entry (a crashed session, a name the registry has not yet released) must not block a
* legitimate relaunch. Each attempt also carries its own {@link #startAwaitingShellPrompt} retry.
*
* @param nameForAttempt called once per attempt (0-based) to produce that attempt's name
* @throws HerdrException the last {@code agent_name_taken} failure once {@code nameRetries} is
* spent, or immediately for any other herdr failure
*/
public static Agent startUniquelyNamed(AgentControl agents, String kind, List<String> args,
String paneId, IntFunction<String> nameForAttempt,
int nameRetries, int shellReadyRetries, Runnable sleeper) {
HerdrException last = null;
for (int attempt = 0; attempt < nameRetries; attempt++) {
String name = nameForAttempt.apply(attempt);
try {
return startAwaitingShellPrompt(agents, name, kind, args, paneId,
shellReadyRetries, sleeper);
} catch (HerdrException e) {
if (!"agent_name_taken".equals(e.code())) throw e;
log.debug("agent name '{}' taken, retrying", name);
last = e;
}
}
throw last;
}
}
@@ -5,6 +5,7 @@ import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.PendingCloseMarker;
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -12,6 +13,7 @@ import dev.ltms.fleet.launch.ClaudeCodeArguments;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.security.SecureRandom;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
@@ -19,6 +21,7 @@ import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.atomic.AtomicLong;
import dev.ltms.fleet.peer.PeerLauncher;
/**
@@ -83,6 +86,13 @@ public final class LeadLauncher {
private final AgentControl agents;
private final WorkspaceControl spaces;
private final FleetConfig cfg;
private final Runnable sleeper;
// Per-process token mixed into each lead agent name so a fresh daemon process (seq back at 0)
// cannot collide with a same-name lead that outlived a restart — the same scheme
// HerdrPeerLauncher uses for members (fleetd #727).
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
private final AtomicLong nameSeq = new AtomicLong();
/**
* @param agents herdr agent control (start, list)
@@ -90,9 +100,27 @@ public final class LeadLauncher {
* @param cfg the loaded config — {@code fleet.leaders}, {@code profiles} and each lead's tab
*/
public LeadLauncher(AgentControl agents, WorkspaceControl spaces, FleetConfig cfg) {
this(agents, spaces, cfg, () -> sleepUninterruptibly(300));
}
/**
* Test seam: as above, plus an injectable {@code sleeper} for the {@code agent_pane_busy}
* retry (fleetd #727), so a test can prove the retry budget without a real sleep.
*/
LeadLauncher(AgentControl agents, WorkspaceControl spaces, FleetConfig cfg, Runnable sleeper) {
this.agents = agents;
this.spaces = spaces;
this.cfg = cfg;
this.sleeper = sleeper;
}
/** Uninterruptible sleep — the production {@link #sleeper} between {@code agent_pane_busy} retries. */
private static void sleepUninterruptibly(long ms) {
try {
Thread.sleep(ms);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
/**
@@ -306,7 +334,16 @@ public final class LeadLauncher {
return null;
}
/** Start one lead. Returns false (having logged) rather than throwing on any failure. */
/**
* Start one lead. Returns false (having logged) rather than throwing on any failure.
*
* <p>Goes through the same {@link ResilientAgentLaunch} seam every member spawn uses
* (fleetd #727): the assembled argv is refused outright if it cannot fit the pane line herdr
* types it into, a stale {@code agent_name_taken} (a crashed session's name the registry has
* not yet released) is retried under a fresh per-attempt name rather than refusing the whole
* relaunch, and a seed pane whose shell has not reached its prompt yet ({@code
* agent_pane_busy}) is retried rather than failing on the first miss.
*/
private boolean launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
String label = lead.tabLabel();
String cwd = (lead.cwd() == null || lead.cwd().isBlank())
@@ -323,8 +360,12 @@ public final class LeadLauncher {
// Same shape as the member launchers: herdr resolves the executable from `kind`, so
// argv[0] (the configured launcher, e.g. `ccs`) is dropped and only the rest is passed.
List<String> argv = leadArgv(profile);
Agent started = agents.start("lead-" + name, herdrKind(profile),
argv.isEmpty() ? argv : argv.subList(1, argv.size()), tab.rootPaneId());
ResilientAgentLaunch.checkFits(profile.profile(), argv);
List<String> args = argv.isEmpty() ? argv : argv.subList(1, argv.size());
Agent started = ResilientAgentLaunch.startUniquelyNamed(agents, herdrKind(profile), args,
tab.rootPaneId(),
attempt -> "lead-" + name + "-" + nameNonce + "-" + nameSeq.incrementAndGet(),
ResilientAgentLaunch.NAME_RETRIES, ResilientAgentLaunch.SHELL_READY_RETRIES, sleeper);
// Label AFTER the start succeeds. A label written before would survive a failed start
// and then read back as a live lead on the next boot, which is the exact staleness the
@@ -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,
@@ -6,6 +6,7 @@ import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -67,16 +68,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
private static final Logger log = LoggerFactory.getLogger(HerdrPeerLauncher.class);
/** herdr rejects a duplicate agent {@code name}; we retry a bumped name this many times. */
private static final int NAME_RETRIES = 8;
/**
* Retries for {@code agent.start} against a seed pane whose shell has not reached its prompt
* yet — {@code tab.create}/{@code pane.split} return as soon as the pane exists, and herdr
* refuses to start an agent in a pane that is not "an available shell" ({@code agent_pane_busy}).
*/
private static final int SHELL_READY_RETRIES = 20;
private final String namePrefix; // label prefix: naming + reap scheme
private final AgentControl agents;
private final WorkspaceControl spaces;
@@ -766,90 +757,26 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
// Protocol 19 resolves the executable from the agent kind (== namePrefix here), so
// argv[0] — the configured executable — is dropped and only the extra args are passed.
List<String> args = argv.isEmpty() ? argv : argv.subList(1, argv.size());
checkPaneCommandFits(cfg, argv);
HerdrException last = null;
for (int attempt = 0; attempt < NAME_RETRIES; attempt++) {
long seq = nameSeq.incrementAndGet();
String name = namePrefix + "-" + cfg.profile() + "-" + nameNonce + "-" + seq;
try {
return new Started(startAwaitingShellPrompt(name, args, paneId), seq);
} catch (HerdrException e) {
if (!"agent_name_taken".equals(e.code())) throw e;
log.debug("peer name '{}' taken, retrying", name);
last = e;
}
try {
ResilientAgentLaunch.checkFits(cfg.profile(), argv);
} catch (ResilientAgentLaunch.TooLargeException e) {
throw new PeerUnreachableException(e.getMessage());
}
throw last;
}
/**
* fleetd #220: herdr does not exec the launch command — it TYPES it into the pane as one line,
* and a pty line buffer holds only {@value #PANE_COMMAND_BYTE_LIMIT} bytes (BSD/macOS {@code
* MAX_CANON}). Everything past that byte is dropped. Nothing reports it: herdr answers "agent
* started", the backend exits on the mangled argument it was handed, the pane closes, and the
* only symptom is {@link #waitUntilInjectableOrThrow} timing out 20 seconds later with no
* reason. That is exactly how #214 broke every claude-code spawn — one 50-byte flag pushed a
* 978-byte command to 1028, and the tail that got cut was {@code --autocompact 250000}.
*
* <p>So measure it here and refuse, loudly and immediately, rather than spawn something that
* cannot work. The estimate is deliberately conservative: fleetd cannot see herdr's quoting, so
* every argument is charged its own bytes plus a separator and a quote pair. An over-estimate
* costs a clear error at a length that was already unsafe; an under-estimate would let the
* silent truncation back in.
*
* @throws PeerUnreachableException when the command cannot fit — the same failure the spawn
* would have hit anyway, named at the point it is still
* explainable
*/
private void checkPaneCommandFits(FleetConfig.Profile cfg, List<String> argv) {
int bytes = 0;
String longest = null;
int longestBytes = 0;
for (String arg : argv) {
int argBytes = arg == null ? 0 : arg.getBytes(java.nio.charset.StandardCharsets.UTF_8).length;
bytes += argBytes + QUOTING_OVERHEAD_PER_ARG;
if (argBytes > longestBytes) {
longestBytes = argBytes;
longest = arg;
}
}
if (bytes <= PANE_COMMAND_BYTE_LIMIT) {
return;
}
String culprit = longest == null ? "<none>"
: longest.substring(0, Math.min(longest.length(), 60)) + (longest.length() > 60 ? "…" : "");
throw new PeerUnreachableException(
"launch command for profile " + cfg.profile() + " is about " + bytes + " bytes, over the "
+ PANE_COMMAND_BYTE_LIMIT + "-byte limit of the pane line herdr types it into. "
+ "The pty would drop the tail silently and the backend would exit on a mangled "
+ "argument. Longest argument is " + longestBytes + " bytes: " + culprit
+ " — move it off the command line (a file flag) or shorten it.");
long[] lastSeq = {0};
Agent agent = ResilientAgentLaunch.startUniquelyNamed(agents, namePrefix, args, paneId,
attempt -> {
lastSeq[0] = nameSeq.incrementAndGet();
return namePrefix + "-" + cfg.profile() + "-" + nameNonce + "-" + lastSeq[0];
},
ResilientAgentLaunch.NAME_RETRIES, ResilientAgentLaunch.SHELL_READY_RETRIES, sleeper);
return new Started(agent, lastSeq[0]);
}
/**
* The pty line buffer herdr types a launch command into: BSD/macOS {@code MAX_CANON}. Not a
* fleetd choice and not configurable — see {@link #checkPaneCommandFits}.
* fleetd choice and not configurable — see {@link ResilientAgentLaunch#checkFits}.
*/
static final int PANE_COMMAND_BYTE_LIMIT = 1024;
/** Per-argument allowance for the separating space and a shell quote pair fleetd cannot see. */
private static final int QUOTING_OVERHEAD_PER_ARG = 3;
/** Start the agent into {@code paneId}, waiting out the seed shell's boot with the sleeper. */
private Agent startAwaitingShellPrompt(String name, List<String> args, String paneId) {
HerdrException busy = null;
for (int attempt = 0; attempt < SHELL_READY_RETRIES; attempt++) {
try {
return agents.start(name, namePrefix, args, paneId);
} catch (HerdrException e) {
if (!"agent_pane_busy".equals(e.code())) throw e;
log.debug("pane {} not at its shell prompt yet, retrying agent.start", paneId);
busy = e;
sleeper.run();
}
}
throw busy;
}
static final int PANE_COMMAND_BYTE_LIMIT = ResilientAgentLaunch.PANE_COMMAND_BYTE_LIMIT;
// --- discovery + reap ----------------------------------------------------------------------
@@ -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). */
@@ -349,6 +355,14 @@ public final class MessageService {
*/
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
/**
* Minted once per {@code MessageService} instance and folded into every ticket id (see
* {@link #sendAsync(String, String, Runnable, String)}). {@link #ticketSeq} alone restarts at
* zero for every instance, so without this a ticket id can be reused across instances and
* resolve to an unrelated {@link Task} with no error; this nonce makes that impossible, because
* an id minted by one instance can never match the id space of another.
*/
private final String ticketBootNonce = UUID.randomUUID().toString().substring(0, 6);
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual().name("bridge-async-", 0).factory());
@@ -665,7 +679,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 +938,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 +959,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 +985,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 +1164,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
@@ -1300,7 +1329,7 @@ public final class MessageService {
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
String ticket = "task-" + ticketSeq.incrementAndGet();
String ticket = "task-" + ticketBootNonce + "-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
tasks.put(ticket, task);
if (pushLoop != null) {
@@ -1332,7 +1361,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.
@@ -1,5 +1,6 @@
package dev.ltms.fleet.msg;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
@@ -60,7 +61,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,11 +71,45 @@ 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<>();
private final AtomicLong askSeq = new AtomicLong();
/**
* Minted once per {@code Rendezvous} instance and folded into every {@code turnId} (see
* {@link #openAsk(String)}). {@link #askSeq} alone restarts at zero for every instance, so
* without this a {@code turnId} minted by one instance could be minted again by another and
* resolve to an unrelated ask with no error; this nonce makes that impossible, because an id
* minted by one instance can never match the id space of another.
*/
private final String askBootNonce = UUID.randomUUID().toString().substring(0, 6);
/** Per-session index of the currently-open ask, so duplicate fleet_ask calls coalesce onto one turn. */
private final ConcurrentHashMap<String, String> openAsksBySession = new ConcurrentHashMap<>();
@@ -89,8 +124,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 +149,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 +169,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();
}
/**
@@ -145,9 +196,9 @@ public final class Rendezvous {
while (true) {
AskWaiter[] minted = { null };
String turnId = openAsksBySession.computeIfAbsent(session, _ -> {
String newTurnId = session + "#" + askSeq.incrementAndGet();
String newTurnId = session + "#" + askBootNonce + "-" + 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 +234,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 +306,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}) {
@@ -1,10 +1,16 @@
package dev.ltms.fleet.lead;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
import dev.ltms.fleet.herdr.WorkspaceControl;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.LinkedHashMap;
import java.util.List;
@@ -64,6 +70,21 @@ class LeadLauncherTest {
return new LeadLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), cfg);
}
/** As {@link #launcher}, plus a fast no-op sleeper so a busy-retry test never real-sleeps. */
private static LeadLauncher fastLauncher(FakeHerdr herdr, FleetConfig cfg) {
return new LeadLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), cfg, () -> { });
}
/** {@link #opusProfile()} with one argv element long enough to overflow the pane line limit. */
private static FleetConfig.Profile hugeArgvProfile() {
return new FleetConfig.Profile(
"opus", null, "claude-opus-5", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "x".repeat(1500)), "tab", "fleet", null,
"http://127.0.0.1:8765/mcp", null, null,
null, null, null,
Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true, null);
}
@SuppressWarnings("unchecked")
private static List<String> startedArgs(FakeHerdr herdr) {
return (List<String>) ((Map<String, Object>) herdr.lastCall("agent.start").params()).get("args");
@@ -86,7 +107,11 @@ class LeadLauncherTest {
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertTrue(herdr.called("agent.start"), "a lead must actually be started");
assertEquals("lead-opus", startedName(herdr));
// fleetd #727: the name carries a per-process nonce and a per-start sequence number — the
// same unique-naming scheme HerdrPeerLauncher uses for members — rather than the fixed
// "lead-opus" a stale registry entry could block a legitimate relaunch under.
assertTrue(startedName(herdr).matches("lead-opus-[0-9a-f]{6}-\\d+"),
"name is lead-<name>-<nonce>-<seq>: " + startedName(herdr));
}
/** The tab is labelled with the configured `tab:` so the scanner finds the lead on the next resolve. */
@@ -436,4 +461,114 @@ class LeadLauncherTest {
assertEquals(0, launcher(herdr, cfg).ensureLeads());
assertTrue(herdr.calls.isEmpty(), "nothing declared ⇒ nothing scanned");
}
// ── fleetd #727: the same three launch protections every member spawn gets ───────────────────
/**
* A freshly created pane may not have redrawn its prompt yet, so herdr answers
* {@code agent_pane_busy}. The lead launch must wait it out rather than fail on the first miss
* — exactly the retry {@code HerdrPeerLauncher} already gives every member.
*/
@Test
void aLeadLaunchRetriesWhileTheSeedShellBoots() {
FakeHerdr herdr = new FakeHerdr().agentPaneBusyTimes(2);
assertEquals(1, fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"the lead must still start once the shell is ready");
assertEquals(3, herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"two busy rejections, then the successful start");
}
/**
* The busy retry is bounded, not an infinite poll. If the pane never becomes ready the launch
* must eventually give up and log the failure, not hang the daemon's reconcile loop forever —
* proven here by a budget that would still be busy on attempt
* {@value ResilientAgentLaunch#SHELL_READY_RETRIES} and a call count that stops exactly there.
*/
@Test
void aLeadLaunchGivesUpAfterTheBoundedBusyBudgetRatherThanLoopingForever() {
FakeHerdr herdr = new FakeHerdr().agentPaneBusyTimes(999);
assertEquals(0, fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"a pane that never becomes ready must not be reported as a started lead");
assertEquals(ResilientAgentLaunch.SHELL_READY_RETRIES,
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"the retry budget is bounded: it stops after exactly SHELL_READY_RETRIES attempts");
}
/**
* herdr refuses a duplicate agent {@code name} outright. A stale registry entry — a crashed
* lead session, or a name the registry has not yet released — must not permanently block a
* legitimate relaunch, so each retry attempt carries a fresh per-start name.
*/
@Test
@SuppressWarnings("unchecked")
void aLeadLaunchRetriesUnderAFreshNameWhenTheOldNameIsStillTaken() {
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(2);
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"the lead must still start once a free name is found");
List<String> names = herdr.calls.stream()
.filter(c -> c.method().equals("agent.start"))
.map(c -> ((Map<String, Object>) c.params()).get("name").toString())
.toList();
assertEquals(3, names.size(), "2 rejected + 1 success");
assertEquals(3, Set.copyOf(names).size(), "each attempt must use a distinct name");
}
/**
* The name-collision retry is bounded too. If the name is taken on every attempt, the launch
* must give up rather than keep minting new names forever — the per-start naming scheme means
* every failed attempt was refused outright by herdr (no process exists under a name herdr
* refused), so a bounded, exhausted retry never leaves a second process running: nothing is
* running at all.
*/
@Test
void aLeadLaunchNeverEndsUpWithASecondProcessWhenTheNameStaysTaken() {
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(999);
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"a name that is never free must not be reported as a started lead");
assertEquals(ResilientAgentLaunch.NAME_RETRIES,
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"the retry budget is bounded: it stops after exactly NAME_RETRIES attempts, "
+ "never racing a duplicate into existence");
}
/**
* fleetd #220/#727: herdr types the launch command into the pane as one line, and a pty line
* buffer holds only 1024 bytes — past that the tail is dropped with no error at all, and the
* backend exits on a mangled argument. The lead launch must refuse an over-long command outright
* rather than let it be typed and silently truncated.
*/
@Test
void anOverlongLeadArgvIsRefusedRatherThanTypedAndTruncated() {
FakeHerdr herdr = new FakeHerdr();
Logger logger = (Logger) LoggerFactory.getLogger(LeadLauncher.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
int started;
try {
started = launcher(herdr, configWith(lead("opus", "lead: opus", 1), hugeArgvProfile()))
.ensureLeads();
} finally {
logger.detachAppender(appender);
}
assertEquals(0, started, "an over-long command must never be reported as a started lead");
assertFalse(herdr.called("agent.start"),
"nothing may be started — a truncated command is worse than no spawn");
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.contains("failed to launch"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a WARN naming the launch failure: "
+ appender.list));
assertTrue(warn.contains("1024"), "names the limit: " + warn);
assertTrue(warn.contains("opus"), "names the profile: " + warn);
assertTrue(warn.contains("x".repeat(60)), "names the culprit argument: " + warn);
}
}
@@ -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");
@@ -888,6 +888,39 @@ class MessageServiceTest {
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
}
// --- fleetd #719: a per-boot nonce keeps one instance's ticket ids out of another's space ---
/** A second, fully independent instance — its own agents/injector/rendezvous/inbox, not shared. */
private MessageService newIndependentInstance() {
FakeHerdr otherHerdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
AgentControl otherAgents = new AgentControl(otherHerdr);
Injector otherInjector = new Injector(otherAgents);
return new MessageService(otherAgents, otherInjector, new Rendezvous(), new InMemoryReplyInbox());
}
@Test
void twoInstancesMintDisjointTicketIds() {
MessageService other = newIndependentInstance();
String ticketFromThis = messages.sendAsync(T, "task on first instance", null, null);
String ticketFromOther = other.sendAsync(T, "task on second instance", null, null);
assertNotEquals(ticketFromThis, ticketFromOther,
"each instance mints its own id space, so even a first ticket from each must differ");
}
@Test
void foreignInstanceTicketDoesNotResolve() {
MessageService other = newIndependentInstance();
String ticket = messages.sendAsync(T, "task on first instance", null, null);
// `other` must reach the same sequence number, or this test passes against an empty map
// instead of against a colliding id.
other.sendAsync(T, "task on second instance", null, null);
// control: the id resolves in the instance that minted it, so a null below cannot be
// explained by broken plumbing — only by the ticket being foreign to `other`.
assertNotNull(messages.poll(ticket), "the minting instance must still resolve its own ticket");
assertNull(other.poll(ticket), "a ticket minted by a different instance must not resolve here");
}
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
@Test
@@ -967,11 +1000,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 +1034,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 +1061,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 +1072,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 +1093,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 +1108,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 +1128,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 +1256,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 +1276,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 +1407,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 +1447,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 +1484,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 +1513,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 +1567,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 +1673,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 +1720,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 +1786,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 +1820,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 +1841,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 +1941,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 +1976,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 +2004,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 +2335,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 +2361,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 +2754,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 +2766,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 +2783,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,127 @@ 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");
}
// ── fleetd #729: per-boot nonce guards turnId against cross-instance reuse ────────────────
@Test
void twoInstancesMintDisjointTurnIds() {
Rendezvous other = new Rendezvous();
Rendezvous.AskTicket fromThis = rendezvous.openAsk(W);
Rendezvous.AskTicket fromOther = other.openAsk(W);
assertNotEquals(fromThis.turnId(), fromOther.turnId(),
"each instance mints its own id space, so even a first ask from each must differ");
}
@Test
void foreignInstanceTurnIdDoesNotResolve() {
Rendezvous other = new Rendezvous();
// `other` must reach the same sequence number as `rendezvous` (two asks each, the first
// closed so the second mints fresh), or this test passes against an empty map instead of
// against a colliding id.
Rendezvous.AskTicket firstFromThis = rendezvous.openAsk(W);
rendezvous.closeAsk(firstFromThis.turnId());
Rendezvous.AskTicket secondFromThis = rendezvous.openAsk(W);
Rendezvous.AskTicket firstFromOther = other.openAsk(W);
other.closeAsk(firstFromOther.turnId());
other.openAsk(W);
// control: the id resolves in the instance that minted it, so a false below cannot be
// explained by broken plumbing — only by the turnId being foreign to `other`.
assertTrue(rendezvous.answerAsk(secondFromThis.turnId(), "answer from this instance"),
"the minting instance must still resolve its own turnId");
assertFalse(other.answerAsk(secondFromThis.turnId(), "answer from other instance"),
"a turnId minted by a different instance must not resolve here");
}
@Test
void openAskStillCoalescesDuplicatesAndStillMintsDistinctIdsPerAsk() {
Rendezvous.AskTicket t1 = rendezvous.openAsk(W);
Rendezvous.AskTicket t2 = rendezvous.openAsk(W);
assertEquals(t1.turnId(), t2.turnId(),
"a second openAsk while one is open still coalesces onto the same turn");
assertFalse(t2.fresh(), "the coalesced ask is still reported as not fresh");
rendezvous.closeAsk(t1.turnId());
Rendezvous.AskTicket t3 = rendezvous.openAsk(W);
assertNotEquals(t1.turnId(), t3.turnId(), "two asks from the same session still get different turnIds");
}
@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());
}
}