Compare commits
15 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1fc9e85bf1 | |||
| 809b7d9b20 | |||
| 8cf7215d56 | |||
| 1ef93e57cc | |||
| cf0c9b9316 | |||
| 337dbd491e | |||
| 7cf6075b79 | |||
| fbdcd709c9 | |||
| 8d3f10d291 | |||
| 38f4fd64ee | |||
| 3fc39b981d | |||
| 70328ca0f8 | |||
| 2374de28e4 | |||
| dd18bd1f38 | |||
| ed4f4b08ad |
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user