Compare commits
26 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6cccd458d4 | |||
| a639969a9a | |||
| 17c3a69c57 | |||
| 49a5875586 | |||
| 634d33b50b | |||
| 1db79bcaa9 | |||
| 4507bc5a70 | |||
| d7239ed23b | |||
| dfeb9340b4 | |||
| 1513d4f260 | |||
| 4ca7d72303 | |||
| c0545d003d | |||
| 275ac0d251 | |||
| c1e06c9e12 | |||
| 204da67d66 | |||
| b091c51eee | |||
| 84d631b030 | |||
| 3f8c38fc54 | |||
| ed2fd6646a | |||
| 5441a2b321 | |||
| 384867dfa3 | |||
| f40c19ecf0 | |||
| 034e17bb32 | |||
| 68b428c484 | |||
| e20ccab1eb | |||
| d83821bbce |
@@ -139,12 +139,21 @@ prefer `wait:false` + `fleet_poll` for anything non-trivial: a blocking `fleet_s
|
||||
**Delegating does not delegate responsibility.** Workers open PRs; you are the gate. Never delegate
|
||||
the merge — and merging on a reviewer's word is delegating it by proxy.
|
||||
|
||||
**When a decision blocks you, consult architects — not the operator.** Spawn one or more architect
|
||||
members, give them the question and the evidence you have, and act on what they agree. They are
|
||||
authorized to settle it, not only to advise. If two of them still disagree after two rounds, they
|
||||
return both positions and you decide. Go to the operator only for something outside the fleet's
|
||||
authority: money, credentials, or a promise made to someone else. **Then write the decision on the
|
||||
ticket.** Taking the operator out of the loop also removes the signal they used to get, because
|
||||
that signal was the block itself — work stopped, so they found out. A ticket comment replaces it,
|
||||
and it reaches them whether or not they are at a terminal when you decide.
|
||||
|
||||
| Intent | Tool |
|
||||
|---|---|
|
||||
| Confirm your own role | `fleet_whoami` |
|
||||
| See backends available | `fleet_profiles` |
|
||||
| Start a member | `fleet_spawn{role?, profile?, cwd?, worktree?, ticket?, sessionName?, resumeSessionId?}` → `sessionId` + `paneId` |
|
||||
| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) · one peer's state: `fleet_status{sessionId}` |
|
||||
| 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` |
|
||||
|
||||
@@ -5,6 +5,21 @@
|
||||
here took a revert and two upstream fixes — see §7.1, which is the useful part of this document. One
|
||||
risk is **accepted rather than solved**: a stream cut by any mid-response timer arrives as HTTP 200
|
||||
with no terminator, and our third-party members cannot detect it (§7.2).
|
||||
|
||||
> **Superseded in part — 2026-09-13.** Two claims on this page are no longer true of the live fleet.
|
||||
> I measured both on this host today.
|
||||
>
|
||||
> 1. **The model is named `acoder` now, not `deepseek-v4-flash`.** `acoder` is a stable alias, and
|
||||
> the model behind it changed on 2026-08-28: it is Qwen3.8-27B, not DeepSeek. The old name is
|
||||
> still served, so nothing broke — the gateway answers it and reports `"model": "acoder"` in the
|
||||
> reply, which is how you can see for yourself that it is an alias. `fleetd.yaml` moved to
|
||||
> `acoder` on 2026-09-13. Do not guess behaviour from the name; ask the gateway's own manifest,
|
||||
> `GET https://llm.ltms.dev/v1/deployment`, and read its `generation` field.
|
||||
> 2. **`local` sits at `weight: 0`, not 100.** Only `gx` is auto-selected today.
|
||||
>
|
||||
> §2 and §3 below are the plan as written in August. They are the record of the migration, so they
|
||||
> stay as they are. If this note stops matching `fleetd.yaml`, re-measure and rewrite the note.
|
||||
|
||||
· **Upstream:** [systems/vms wiki → LLM and MCP Gateway](https://git.ltms.dev/systems/vms/wiki/LLM-and-MCP-Gateway)
|
||||
· **Upstream issue:** [systems/vms#31](https://git.ltms.dev/systems/vms/issues/31)
|
||||
|
||||
@@ -380,7 +395,10 @@ one turn; this one costs the whole task and is indistinguishable from a slow wor
|
||||
|
||||
- Token accepted on both surfaces. **Unauthenticated → 401**, so the Caddy proxy really does gate —
|
||||
the wiki's "SecurityPolicy fails open" warning is about the gateway itself, not the edge.
|
||||
- `/v1/models` returns exactly `["deepseek-v4-flash"]`, so trap 3 is clear.
|
||||
- `/v1/models` returned exactly `["deepseek-v4-flash"]` **on 2026-08-15**, so trap 3 was clear
|
||||
then. It returns 6 ids now — `acoder`, `qwen3.8-27b-nvfp4`, `deepseek-v4-flash` and three
|
||||
embedding names — measured on this host 2026-09-13. The exact-name rule still holds; the
|
||||
one-item list does not.
|
||||
- **Reasoning survives both surfaces** — see §3b above.
|
||||
- The launcher's generated opencode provider block is correct, carrying a real 48-character `llmk-`
|
||||
key rather than the `fleetd-local-noauth` placeholder.
|
||||
|
||||
@@ -21,8 +21,10 @@ import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.LoopWatchdog;
|
||||
import dev.ltms.fleet.inject.StatusPoller;
|
||||
import dev.ltms.fleet.inject.TurnListener;
|
||||
import dev.ltms.fleet.inject.TurnRegistrar;
|
||||
import dev.ltms.fleet.inject.MemberPresence;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
@@ -79,6 +81,9 @@ import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.BooleanSupplier;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Predicate;
|
||||
@@ -491,61 +496,34 @@ public final class Fleetd {
|
||||
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
||||
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
|
||||
MemberPresence presence = sessions.asPresence();
|
||||
TurnListener turnListener = new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
completion.onTurnComplete(target);
|
||||
sessions.onTurnComplete(target);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasPostTurnAction(String target) {
|
||||
return sessions.hasPostTurnAction(target);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean onTurnCompleteWithPostAction(String target) {
|
||||
completion.resolveBeforePostAction(target);
|
||||
return sessions.onTurnCompleteWithPostAction(target);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
|
||||
completion.onDelivered(target, token);
|
||||
sessions.onDelivered(target, token);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnFailed(String target) {
|
||||
completion.onTurnFailed(target);
|
||||
sessions.onTurnFailed(target);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnFailed(String target, String reason) {
|
||||
completion.onTurnFailed(target, reason);
|
||||
sessions.onTurnFailed(target);
|
||||
}
|
||||
};
|
||||
// fleetd #561: extracted to a static factory (see turnListener below) — the anonymous class
|
||||
// this replaced had two bare, unguarded statements per callback, and nothing enforced that
|
||||
// the completion half went first beyond call order in the source.
|
||||
TurnListener turnListener = turnListener(completion, sessions);
|
||||
Predicate<String> deliverable = deliverableTo(presence, leads);
|
||||
// fleetd #556: registration is wired directly to `completion`, not folded into the
|
||||
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
|
||||
// listener) throwing, regardless of call order. See TurnRegistrar's javadoc.
|
||||
// fleetd #589 Group 3 (:505): extracted to turnRegistrar(...) — see FleetdTurnRegistrarWiringTest.
|
||||
Injector injector = new Injector(router, turnListener, deliverable,
|
||||
presence::forget, completion::register);
|
||||
presence::forget, turnRegistrar(completion));
|
||||
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||
poller.start();
|
||||
|
||||
// CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or
|
||||
// unusable), fleetd stays soft-state on the in-memory inbox. The AMQP inbox owns a broker
|
||||
// connection, so keep the reference to close it in the ordered shutdown hook.
|
||||
final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::open);
|
||||
// fleetd #589 Group 3 (:512): opener extracted to replyInboxOpener() — see
|
||||
// FleetdReplyInboxOpenerWiringTest.
|
||||
final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), replyInboxOpener());
|
||||
// CB-637: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate
|
||||
// broker from the reply inbox by design (see FleetConfig.Coordinator). Absent a coordinator:
|
||||
// block this is null and every lead path below is simply not wired, which is exactly the
|
||||
// behaviour before this ticket. It owns a broker connection, so keep the reference for the
|
||||
// ordered shutdown hook.
|
||||
final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), LeadMailbox::open);
|
||||
// fleetd #589 Group 3 (:518): opener extracted to leadMailboxOpener() — see
|
||||
// FleetdLeadMailboxOpenerWiringTest.
|
||||
final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), leadMailboxOpener());
|
||||
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
|
||||
// The pin also feeds CallerResolver below: a primary running inside a herdr pane would
|
||||
// otherwise resolve as a worker and be refused every orchestration tool.
|
||||
@@ -615,10 +593,12 @@ public final class Fleetd {
|
||||
// fleetd #386: System::nanoTime freezes across a macOS sleep, so the stall check also
|
||||
// gets a wall-clock source to detect and correct for that freeze. Every other decision
|
||||
// in FleetHealthMonitor stays on the monotonic clock, unchanged.
|
||||
// fleetd #589 Group 3 (:591): failTarget extracted to healthFailTarget(...) — see
|
||||
// FleetdHealthFailTargetWiringTest.
|
||||
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
|
||||
System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()),
|
||||
cfg.health().intervalOrDefault(),
|
||||
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
|
||||
cfg.health().workingSuspectAfterOrDefault(), healthFailTarget(messages));
|
||||
String coverage = FleetHealthMonitor.coverage(true,
|
||||
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
||||
if ("detection-only".equals(coverage)) {
|
||||
@@ -638,27 +618,11 @@ public final class Fleetd {
|
||||
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
|
||||
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
|
||||
// reached /metrics — the delegation was unresolvable and nothing said so.
|
||||
sessions.onRelease(detail -> {
|
||||
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
|
||||
// where to re-dispatch onto the same tree, not just that the worker vanished.
|
||||
String reason = "the worker session was released before it replied";
|
||||
if (detail.worktreePath() != null) {
|
||||
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
||||
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
||||
}
|
||||
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
||||
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
||||
if (detail.agentSessionId() != null) {
|
||||
reason += " agentSessionId=" + detail.agentSessionId();
|
||||
}
|
||||
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
|
||||
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
|
||||
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
|
||||
// (see MessageService.abandon's javadoc for why those two must differ).
|
||||
messages.abandon(detail.terminalId(), reason, true);
|
||||
replyInbox.release(detail.terminalId());
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
});
|
||||
// fleetd #589 Group 3 (:611-631): the whole cleanup lambda extracted to releaseCleanup(...)
|
||||
// — see FleetdReleaseCleanupWiringTest, and MessageService.abandon's javadoc for the
|
||||
// documented incident (a torn-down worker's rendezvous waiter left open) this lambda exists
|
||||
// to prevent.
|
||||
sessions.onRelease(releaseCleanup(messages, replyInbox, primaryRegistry));
|
||||
|
||||
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
|
||||
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
||||
@@ -696,10 +660,12 @@ public final class Fleetd {
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy);
|
||||
|
||||
FleetMcp.LoopHealthSource loopHealth = loopHealthSource(poller, reaper);
|
||||
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
|
||||
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
|
||||
capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
|
||||
healthCoverageSource(config),
|
||||
loopHealth,
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
@@ -791,7 +757,7 @@ public final class Fleetd {
|
||||
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
|
||||
callers, metrics, deliverable,
|
||||
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
|
||||
quarantineSource, outageSource).build();
|
||||
quarantineSource, outageSource, loopHealth).build();
|
||||
app.start(cfg.bind().host(), cfg.bind().port());
|
||||
log.info("fleetd listening on {}:{}, herdr socket {}",
|
||||
cfg.bind().host(), cfg.bind().port(), socket);
|
||||
@@ -1069,6 +1035,129 @@ public final class Fleetd {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #562 follow-up: package-private factory for {@code fleet_list}'s and {@code
|
||||
* /healthz}'s {@code loopHealth} source, extracted out of {@code main} for the same reason
|
||||
* {@link #capacitySource} and {@link #healthCoverageSource} were. Before this ticket the
|
||||
* {@link FleetMcp.LoopHealthSource} was built inline with a bare {@code new}, so there was
|
||||
* nothing a test could call directly — measured: replacing {@code poller::health} with a
|
||||
* constant {@code () -> LoopWatchdog.State.RUNNING} at the call site compiled clean and left
|
||||
* the full suite green, meaning the daemon could report the {@link StatusPoller} as always
|
||||
* {@code RUNNING} even while it was actually stalled. That is a false negative on the exact
|
||||
* signal this ticket exists to surface, and is the mirror of a false positive muting a real
|
||||
* monitoring component — worse, because there is no noise for anyone to notice and then
|
||||
* silence. {@link FleetdLoopHealthSourceWiringTest} calls this factory directly and pins both
|
||||
* halves separately, plus the {@code reaper == null} branch below.
|
||||
*
|
||||
* <p>{@code reaper} may be {@code null} — a {@link SessionReaper} is only constructed when
|
||||
* {@code lifecycle.idleTtlSeconds} is configured (see the {@code reaper} local above) — and
|
||||
* this factory preserves the existing behaviour of reporting {@link LoopWatchdog.State#STOPPED}
|
||||
* in that case, rather than a {@code NullPointerException} on the first {@code fleet_list} or
|
||||
* {@code /healthz} call.
|
||||
*/
|
||||
static FleetMcp.LoopHealthSource loopHealthSource(StatusPoller poller, SessionReaper reaper) {
|
||||
return new FleetMcp.LoopHealthSource(poller::health,
|
||||
() -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :505}): package-private factory for the {@link Injector}'s
|
||||
* {@link TurnRegistrar}, extracted out of {@code main} for the same reason {@link
|
||||
* #loopHealthSource} was — before this ticket {@code completion::register} was an inline
|
||||
* argument to {@code new Injector(...)}, so nothing could pin it directly. Replacing it with
|
||||
* {@link TurnRegistrar#NOOP} compiles clean and leaves every existing test green: {@code
|
||||
* onDelivered}'s own {@code captureBaseline} does the identical {@code inFlight} check-and-put a
|
||||
* moment later on the ordinary path, so the two are indistinguishable once a turn's delivery
|
||||
* finishes normally. The gap {@link TurnRegistrar}'s own javadoc (fleetd #556) exists to close is
|
||||
* a {@code turnListener} callback throwing between the two — {@link
|
||||
* FleetdTurnRegistrarWiringTest} pins that {@code register} itself (not just {@code
|
||||
* captureBaseline}) makes a delivered turn's waiter resolvable.
|
||||
*/
|
||||
static TurnRegistrar turnRegistrar(CompletionResolver completion) {
|
||||
return completion::register;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :591}): package-private factory for {@link
|
||||
* FleetHealthMonitor}'s {@code failTarget} callback, extracted out of {@code main} for the same
|
||||
* reason {@link #loopHealthSource} was. Before this ticket {@code messages::abandon} was an
|
||||
* inline argument to {@code new FleetHealthMonitor(...)}; replacing it with a no-op {@code
|
||||
* BiConsumer} compiles clean and leaves every existing test green, and in production it means a
|
||||
* member found {@code GONE}/{@code NEVER_READY} never fails the ticket waiting on it — the
|
||||
* caller reports {@code PENDING} for the full 30-minute async timeout instead of the immediate,
|
||||
* accurate failure CB-580 exists to give it. {@link FleetdHealthFailTargetWiringTest} pins that
|
||||
* the returned callback actually reaches the real {@link MessageService#abandon}.
|
||||
*/
|
||||
static BiConsumer<String, String> healthFailTarget(MessageService messages) {
|
||||
return messages::abandon;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :611-631}): package-private factory for the whole {@link
|
||||
* SessionManager#onRelease} cleanup callback, extracted out of {@code main} for the same reason
|
||||
* {@link #loopHealthSource} was. Before this ticket this was an inline lambda built directly
|
||||
* inside {@code main}; replacing its body with a no-op {@code detail -> { }} compiles clean and
|
||||
* leaves every existing test green, and in production it is the exact incident {@link
|
||||
* MessageService#abandon}'s own javadoc documents: a torn-down worker's rendezvous waiter is
|
||||
* left open, so a blocking {@code fleet_send} keeps blocking and an async one reports {@code
|
||||
* PENDING} for a hardcoded thirty minutes on every {@code fleet_stop} and every idle-reap.
|
||||
*
|
||||
* <p>{@link FleetdReleaseCleanupWiringTest} pins all three collaborator calls this lambda makes
|
||||
* — {@code messages.abandon}, {@code replyInbox.release}, and {@code
|
||||
* primaryRegistry.forgetDelegation} — each already tested on its own ({@code MessageServiceTest},
|
||||
* {@code PrimaryRegistryTest}), but never before proven to actually be reached from here.
|
||||
*/
|
||||
static Consumer<SessionManager.ReleaseDetail> releaseCleanup(MessageService messages, ReplyInbox replyInbox,
|
||||
PrimaryRegistry primaryRegistry) {
|
||||
return detail -> {
|
||||
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
|
||||
// where to re-dispatch onto the same tree, not just that the worker vanished.
|
||||
String reason = "the worker session was released before it replied";
|
||||
if (detail.worktreePath() != null) {
|
||||
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
||||
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
||||
}
|
||||
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
||||
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
||||
if (detail.agentSessionId() != null) {
|
||||
reason += " agentSessionId=" + detail.agentSessionId();
|
||||
}
|
||||
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
|
||||
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
|
||||
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
|
||||
// (see MessageService.abandon's javadoc for why those two must differ).
|
||||
messages.abandon(detail.terminalId(), reason, true);
|
||||
replyInbox.release(detail.terminalId());
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :512}): package-private factory for {@code selectReplyInbox}'s
|
||||
* production {@link AmqpOpener}, extracted out of {@code main} for the same reason {@link
|
||||
* #loopHealthSource} was. Before this ticket {@code AmqpReplyInbox::open} was an inline argument
|
||||
* to the {@code selectReplyInbox(...)} call; replacing it with {@code (uri, prefetch) -> new
|
||||
* InMemoryReplyInbox()} compiles clean and leaves every existing test green — {@link
|
||||
* FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own injected opener and
|
||||
* never sees what {@code main} actually passes. {@link FleetdReplyInboxOpenerWiringTest} pins
|
||||
* that this returns the real opener by pointing it at a guaranteed-closed local port and
|
||||
* asserting the real network attempt throws — the inert stub never attempts a connection at all.
|
||||
*/
|
||||
static AmqpOpener replyInboxOpener() {
|
||||
return AmqpReplyInbox::open;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :518}): package-private factory for {@code
|
||||
* openLeadMailbox}'s production {@link LeadMailboxOpener}, extracted out of {@code main} for the
|
||||
* same reason {@link #replyInboxOpener} was — same gap, same fix, the lead-coordination mailbox
|
||||
* instead of the reply inbox. {@link FleetdLeadMailboxOpenerWiringTest} pins that this returns
|
||||
* the real opener the same way.
|
||||
*/
|
||||
static LeadMailboxOpener leadMailboxOpener() {
|
||||
return LeadMailbox::open;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
|
||||
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
|
||||
@@ -1096,6 +1185,157 @@ public final class Fleetd {
|
||||
.orElse(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #561: compose the production {@link TurnListener} from its two halves — the
|
||||
* completion resolver (which resolves a blocked {@code fleet_send}'s waiter) and the session
|
||||
* manager (which drives the member's lifecycle state) — hardened so a throw from either
|
||||
* half's callback can never suppress the other half's callback for the same event.
|
||||
*
|
||||
* <p>Before this, each method below was two bare, unguarded statements: whichever ran first
|
||||
* throwing meant the second one never ran at all, and nothing beyond call order in the
|
||||
* source enforced "the completion half goes first". #556 fixed the identical shape for {@code
|
||||
* onDelivered}'s registration by moving it off the fan-out entirely (see {@code
|
||||
* TurnRegistrar}); this fixes the four remaining callbacks — {@code onTurnComplete}, {@code
|
||||
* onTurnCompleteWithPostAction}, and both {@code onTurnFailed} overloads — by hardening the
|
||||
* fan-out itself instead, since none of them can be pulled out of the listener the way
|
||||
* registration was.
|
||||
*
|
||||
* <p>The invariant this composition guarantees: <b>a throwing session-listener must not
|
||||
* prevent the completion resolver from being told the turn ended.</b> The completion half is
|
||||
* always attempted first, and — as a bonus the resolver does not depend on — its own throw
|
||||
* does not stop the session half from running either. Whatever escapes (from one half or
|
||||
* both) is rethrown once both have been attempted, with a second failure recorded via {@link
|
||||
* Throwable#addSuppressed} on the first, so it still reaches {@link StatusPoller}'s {@code
|
||||
* catch (Throwable)} and logs at ERROR. Nothing here swallows a failure to make the two halves
|
||||
* look safe.
|
||||
*
|
||||
* <p>{@code onTurnCompleteWithPostAction} is the one place order is not just fault-tolerance
|
||||
* but a functional requirement: {@code completion.resolveBeforePostAction} must resolve the
|
||||
* scrape before {@code sessions.onTurnCompleteWithPostAction}'s adapter housekeeping can erase
|
||||
* the pane's rendered output (see {@link CompletionResolver#resolveBeforePostAction}). That is
|
||||
* why this composition does not treat the pair symmetrically the way {@link #bothMustRun}
|
||||
* does for the other three callbacks: the mirror case (completion half throws, session half's
|
||||
* return value still observed) is not preserved here — once the completion half's failure
|
||||
* escapes, the session half's return value is discarded, matching how the {@link Injector}
|
||||
* already treats any throw from this callback as "the action did not start" (see the {@code
|
||||
* started} default at its call site, {@code Injector.java} ~line 616).
|
||||
*
|
||||
* <p>Package-private so {@code FleetdTurnListenerCompositionTest} can build this listener
|
||||
* directly from a real {@link CompletionResolver} and a fake {@link TurnListener} standing in
|
||||
* for {@code sessions}, without booting the rest of {@code main}.
|
||||
*/
|
||||
static TurnListener turnListener(CompletionResolver completion, TurnListener sessions) {
|
||||
return new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
bothMustRun(() -> completion.onTurnComplete(target), () -> sessions.onTurnComplete(target));
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasPostTurnAction(String target) {
|
||||
return sessions.hasPostTurnAction(target);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean onTurnCompleteWithPostAction(String target) {
|
||||
return bothMustRunKeepingSecondResult(() -> completion.resolveBeforePostAction(target),
|
||||
() -> sessions.onTurnCompleteWithPostAction(target));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
|
||||
// fleetd #561: left unguarded on purpose, not because registration survives a
|
||||
// throw elsewhere. completion.onDelivered runs CompletionResolver.captureBaseline,
|
||||
// which already wraps its scrape read in its own catch (RuntimeException) and
|
||||
// fails open (baseline = null) — so this call does not realistically throw, and
|
||||
// there is nothing here for bothMustRun to protect.
|
||||
completion.onDelivered(target, token);
|
||||
sessions.onDelivered(target, token);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnFailed(String target) {
|
||||
bothMustRun(() -> completion.onTurnFailed(target), () -> sessions.onTurnFailed(target));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnFailed(String target, String reason) {
|
||||
bothMustRun(() -> completion.onTurnFailed(target, reason), () -> sessions.onTurnFailed(target));
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #561: run two listener-callback halves for one lifecycle event, guaranteeing the
|
||||
* SECOND always runs even when the FIRST throws. Whatever is thrown is rethrown once both
|
||||
* halves have been attempted — a second failure is attached to the first via {@link
|
||||
* Throwable#addSuppressed} rather than dropped. Never swallows.
|
||||
*/
|
||||
private static void bothMustRun(Runnable completionHalf, Runnable sessionsHalf) {
|
||||
Throwable failure = null;
|
||||
try {
|
||||
completionHalf.run();
|
||||
} catch (Throwable t) {
|
||||
failure = t;
|
||||
}
|
||||
try {
|
||||
sessionsHalf.run();
|
||||
} catch (Throwable t) {
|
||||
if (failure == null) {
|
||||
failure = t;
|
||||
} else {
|
||||
failure.addSuppressed(t);
|
||||
}
|
||||
}
|
||||
if (failure != null) {
|
||||
throwUnchecked(failure);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #561: like {@link #bothMustRun}, but for {@code onTurnCompleteWithPostAction}, whose
|
||||
* session half returns the value the {@link Injector} needs. The second (session) half's
|
||||
* result is what this method returns; if the first (completion) half throws, the second half
|
||||
* still runs and its result is still computed here, but the throw is rethrown afterward
|
||||
* regardless — so that result is discarded at the {@link Injector} call site exactly as it
|
||||
* already is today when this callback throws (see {@link #turnListener}'s javadoc).
|
||||
*/
|
||||
private static boolean bothMustRunKeepingSecondResult(Runnable completionHalf,
|
||||
BooleanSupplier sessionsHalf) {
|
||||
Throwable failure = null;
|
||||
try {
|
||||
completionHalf.run();
|
||||
} catch (Throwable t) {
|
||||
failure = t;
|
||||
}
|
||||
boolean result = false;
|
||||
try {
|
||||
result = sessionsHalf.getAsBoolean();
|
||||
} catch (Throwable t) {
|
||||
if (failure == null) {
|
||||
failure = t;
|
||||
} else {
|
||||
failure.addSuppressed(t);
|
||||
}
|
||||
}
|
||||
if (failure != null) {
|
||||
throwUnchecked(failure);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #561: rethrow a captured {@link Throwable} without a checked-exception wrapper. The
|
||||
* two callback halves above never declare a checked exception (both existing production
|
||||
* halves — {@code CompletionResolver} and {@code SessionManager} — only ever throw unchecked),
|
||||
* so this only ever actually rethrows a {@link RuntimeException} or {@link Error}; the generic
|
||||
* cast is the standard "sneaky throw" idiom, not a claim that a checked exception is expected.
|
||||
*/
|
||||
@SuppressWarnings("unchecked")
|
||||
private static <T extends Throwable> void throwUnchecked(Throwable t) throws T {
|
||||
throw (T) t;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #480: construct the {@link LeadRollover} executor only when {@code leadRollover:} is
|
||||
* present at startup — the same presence gate {@code leadHeartbeat:} uses just above this
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.inject;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import org.slf4j.Logger;
|
||||
@@ -46,6 +47,17 @@ import java.util.stream.Collectors;
|
||||
* turn — the pickup-grace path (a turn too fast to sample) unwedges the queue but does not fire
|
||||
* completion, since without a sampled {@code working} there is no trustworthy "the worker just
|
||||
* finished the task" signal to act on.
|
||||
*
|
||||
* <p><strong>Delivery honesty (fleetd #551).</strong> The one external send call
|
||||
* ({@link AgentControl#send}) is irreversible, and its own response can still fail after the text
|
||||
* has already reached the worker's pane — herdr replied (or may have), so the request was already
|
||||
* processed, but the reply itself then failed to parse or carried an error. A queued entry is
|
||||
* polled off the queue and marked {@code ATTEMPTED} <em>before</em> that call is made, not after,
|
||||
* so a failure from the call itself can never be recorded as a confident {@code NOT_DELIVERED} for
|
||||
* text that may already be sitting in the pane. {@code NOT_DELIVERED} stays reserved for the cases
|
||||
* where nothing was ever attempted — the readiness grace expiring, {@link #drop}, or a herdr
|
||||
* {@code *_not_found} error, which the rest of this codebase already treats as a confirmed absence
|
||||
* rather than a merely inconclusive failure (see {@code StatusPoller}, {@code AgentControl}).
|
||||
*/
|
||||
public final class Injector {
|
||||
|
||||
@@ -217,8 +229,23 @@ public final class Injector {
|
||||
|
||||
/** The result of trying to remove an undelivered message from the injector. */
|
||||
public enum Cancellation {
|
||||
/** The message was still queued and this call removed it; the target saw nothing. */
|
||||
CANCELLED,
|
||||
/** The message reached the worker's pane and is confirmed delivered. */
|
||||
DELIVERED,
|
||||
/**
|
||||
* fleetd #551: the injector called {@link AgentControl#send} for this message and that call
|
||||
* threw before its outcome was known — the text may or may not have reached the pane. A
|
||||
* caller reporting this to an operator must say "uncertain", not "definitely not
|
||||
* delivered": reading it as a confident negative invites a resend of text that may already
|
||||
* be sitting in the pane (a double delivery), which is worse than the ambiguity itself.
|
||||
*/
|
||||
ATTEMPTED,
|
||||
/**
|
||||
* The message never reached the worker's pane — nothing was ever attempted for it (the
|
||||
* readiness grace expired, the target was dropped, or send failed with a herdr
|
||||
* {@code *_not_found} error, which is a confirmed absence, not merely inconclusive).
|
||||
*/
|
||||
NOT_DELIVERED
|
||||
}
|
||||
|
||||
@@ -241,7 +268,29 @@ public final class Injector {
|
||||
|
||||
/** A pending message and the future that completes when it has been delivered. */
|
||||
private static final class Pending {
|
||||
enum State { QUEUED, DELIVERED, NOT_DELIVERED, CANCELLED }
|
||||
enum State {
|
||||
/** Still in the target's queue, not yet attempted. */
|
||||
QUEUED,
|
||||
/**
|
||||
* fleetd #551: {@link AgentControl#send} has been called for this entry and its outcome
|
||||
* is not yet known — recorded BEFORE the call (see {@link #onStatus}), so a Throwable
|
||||
* from send() can never leave the entry either still QUEUED or falsely marked
|
||||
* NOT_DELIVERED. Upgraded to DELIVERED on success; left as ATTEMPTED on an ordinary
|
||||
* failure, since reaching the catch does not prove the text never reached the pane.
|
||||
*/
|
||||
ATTEMPTED,
|
||||
/** {@code send()} returned normally: the text is confirmed to have reached the pane. */
|
||||
DELIVERED,
|
||||
/**
|
||||
* Confirmed — not merely inconclusive — that nothing was ever sent for this entry: the
|
||||
* worker never became ready ({@link #onStatus}'s readiness-grace expiry), the target was
|
||||
* dropped ({@link #drop}), or send failed with a herdr {@code *_not_found} error. Never
|
||||
* written for an entry whose send() outcome is unknown; see ATTEMPTED.
|
||||
*/
|
||||
NOT_DELIVERED,
|
||||
/** Removed from the queue by {@link #cancel} before it was ever attempted. */
|
||||
CANCELLED
|
||||
}
|
||||
|
||||
final String target;
|
||||
final String text;
|
||||
@@ -333,7 +382,14 @@ public final class Injector {
|
||||
}
|
||||
|
||||
private static Cancellation cancellationOf(Pending p) {
|
||||
return p.state == Pending.State.DELIVERED ? Cancellation.DELIVERED : Cancellation.NOT_DELIVERED;
|
||||
// fleetd #551 (comment 17037): a two-way split on a now-three-way question folded ATTEMPTED
|
||||
// into NOT_DELIVERED with no compiler error and no failing test — the exact defect this
|
||||
// ticket exists to fix, one layer up. ATTEMPTED gets its own answer instead.
|
||||
return switch (p.state) {
|
||||
case DELIVERED -> Cancellation.DELIVERED;
|
||||
case ATTEMPTED -> Cancellation.ATTEMPTED;
|
||||
case QUEUED, NOT_DELIVERED, CANCELLED -> Cancellation.NOT_DELIVERED;
|
||||
};
|
||||
}
|
||||
|
||||
private static boolean isQuiescent(Target t) {
|
||||
@@ -426,9 +482,18 @@ public final class Injector {
|
||||
if (p != null && ready.test(target)) {
|
||||
t.notReadySincePoll = 0;
|
||||
t.notReadySinceMillis = 0;
|
||||
// fleetd #551: poll and record BEFORE the irreversible send, not after.
|
||||
// The entry comes off the queue and its state is set to ATTEMPTED here,
|
||||
// unconditionally — so a Throwable escaping the send call below (caught or
|
||||
// not) can never leave the entry QUEUED at the head of t.queue (the fleetd
|
||||
// #546 hazard, since peek() alone would let the next onStatus round re-enter
|
||||
// this block and send the same text again), and no path can write a
|
||||
// confident DELIVERED or NOT_DELIVERED before we actually know which one
|
||||
// happened.
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.ATTEMPTED;
|
||||
try {
|
||||
agentsFor(target).send(target, p.text());
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.DELIVERED;
|
||||
t.awaitingPickup = true;
|
||||
t.awaitingCompletion = true;
|
||||
@@ -436,15 +501,24 @@ public final class Injector {
|
||||
t.injectableSincePickup = 0;
|
||||
sent = p;
|
||||
} catch (Throwable e) {
|
||||
// Delivery failed at herdr; drop the poisoned message and surface it
|
||||
// rather than blocking the queue behind it. Catches Throwable, not just
|
||||
// RuntimeException: fleetd #546 — an Error escaping this send (e.g. a
|
||||
// NoClassDefFoundError, see #413) would otherwise leave the entry QUEUED
|
||||
// at the head of t.queue. Line :378 peeks rather than polls, so the next
|
||||
// onStatus round would re-enter this try and send the same text again,
|
||||
// typing the same brief into the member's pane a second time.
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.NOT_DELIVERED;
|
||||
// fleetd #551: leave p.state == ATTEMPTED (recorded above, before the
|
||||
// call) rather than downgrading it to NOT_DELIVERED here — reaching this
|
||||
// catch does not prove the text never reached the pane. Three of the
|
||||
// four HerdrException throw sites in HerdrCodec fire only after herdr
|
||||
// has already replied (so it processed the request), and the fourth (a
|
||||
// transport IOException) leaves it genuinely unknown whether herdr even
|
||||
// received the bytes — see #551 comment 16867. The one exception is a
|
||||
// herdr `*_not_found` error: that family is already read as "definitely
|
||||
// absent, not merely inconclusive" everywhere else in this codebase
|
||||
// (StatusPoller, AgentControl's own retry, WorkspaceControl,
|
||||
// HerdrPeerLauncher, FleetApp, ReplyPushLoop) because it means the
|
||||
// target pane/agent does not exist at all, so nothing could have been
|
||||
// pasted anywhere — #551 keeps the new state consistent with that
|
||||
// existing vocabulary rather than inventing a second one.
|
||||
if (e instanceof HerdrException he && he.code() != null
|
||||
&& he.code().endsWith("_not_found")) {
|
||||
p.state = Pending.State.NOT_DELIVERED;
|
||||
}
|
||||
sent = p;
|
||||
sendError = e;
|
||||
}
|
||||
|
||||
@@ -20,6 +20,7 @@ import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.inject.LoopWatchdog;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
@@ -108,6 +109,7 @@ public final class FleetMcp {
|
||||
private final Metrics metrics; // CB-502: null → auth failures not counted
|
||||
private final CapacitySource capacity;
|
||||
private final HealthCoverageSource healthCoverage;
|
||||
private final LoopHealthSource loopHealth;
|
||||
private final QuarantineSource quarantine;
|
||||
/** fleetd #201 Unit 5: SEPARATE from {@link #quarantine} — see {@link OutageSource}'s doc. */
|
||||
private final OutageSource outage;
|
||||
@@ -136,6 +138,15 @@ public final class FleetMcp {
|
||||
/** Coverage is supplied by the health wiring, not inferred from a missing dependency. */
|
||||
public record HealthCoverageSource(Supplier<String> value) { }
|
||||
|
||||
/** Progress states for fleetd's singleton background loops, read by {@code fleet_list} and {@code /healthz}. */
|
||||
public record LoopHealthSource(Supplier<LoopWatchdog.State> statusPoller,
|
||||
Supplier<LoopWatchdog.State> sessionReaper) {
|
||||
/** Inert source for callers that do not wire the background loops. */
|
||||
public static LoopHealthSource none() {
|
||||
return new LoopHealthSource(() -> LoopWatchdog.State.STOPPED, () -> LoopWatchdog.State.STOPPED);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id
|
||||
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
|
||||
@@ -328,11 +339,22 @@ public final class FleetMcp {
|
||||
* instead of throwing. See {@link #handover}.
|
||||
*/
|
||||
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
||||
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
|
||||
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
|
||||
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
|
||||
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
|
||||
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
|
||||
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
|
||||
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, authorizationMode, metrics,
|
||||
capacity, healthCoverage, LoopHealthSource.none(), quarantine, leadChannel, outage, leadSeats,
|
||||
peers, leadRollover);
|
||||
}
|
||||
|
||||
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
||||
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
|
||||
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth,
|
||||
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
|
||||
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
|
||||
Objects.requireNonNull(callers, "callers");
|
||||
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
|
||||
== AuthorizationMode.ENFORCED;
|
||||
@@ -343,6 +365,7 @@ public final class FleetMcp {
|
||||
this.outage = Objects.requireNonNull(outage, "outage");
|
||||
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
|
||||
this.healthCoverage = healthCoverage;
|
||||
this.loopHealth = Objects.requireNonNull(loopHealth, "loopHealth");
|
||||
this.leadRollover = leadRollover;
|
||||
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
|
||||
this.transport = HttpServletStreamableServerTransportProvider.builder()
|
||||
@@ -474,7 +497,7 @@ public final class FleetMcp {
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
|
||||
if (denied != null) return denied;
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
|
||||
leadSeats, callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
new CoordinationSource(leadChannel, peers),
|
||||
@@ -805,6 +828,12 @@ public final class FleetMcp {
|
||||
+ "answered (turnId stale)");
|
||||
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,
|
||||
// so the message may already be sitting in the pane. Do not invite a blind retry the way
|
||||
// the case above does; a resend on this route can double-deliver the same brief.
|
||||
case TIMED_OUT_UNCONFIRMED -> text("[no reply within " + timeout + "ms — delivery unconfirmed; "
|
||||
+ "the message may already have reached the worker, so a retry risks sending it "
|
||||
+ "twice — poll status before resending]");
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1546,25 +1575,42 @@ public final class FleetMcp {
|
||||
* @param selfTerm the calling pane's terminal id, or blank for a caller with no pane
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions,
|
||||
Map<String, String> leads, String selfTerm) {
|
||||
Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, null, CapacitySource.none(), new HealthCoverageSource(() -> "off"),
|
||||
QuarantineSource.none(), leads, selfTerm);
|
||||
LoopHealthSource.none(), QuarantineSource.none(), leads, selfTerm);
|
||||
}
|
||||
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, leads, selfTerm,
|
||||
CoordinationSource.none());
|
||||
}
|
||||
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm,
|
||||
LoopHealthSource loopHealth, QuarantineSource quarantine,
|
||||
Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, leads, selfTerm,
|
||||
CoordinationSource.none());
|
||||
}
|
||||
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
LoopHealthSource loopHealth, QuarantineSource quarantine,
|
||||
Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
|
||||
OutageSource.none(), LeadSeatSource.none(), leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none());
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none(), false);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1580,8 +1626,8 @@ public final class FleetMcp {
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination);
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1589,8 +1635,8 @@ public final class FleetMcp {
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination);
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1611,7 +1657,7 @@ public final class FleetMcp {
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
leadSeats, leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
@@ -1631,8 +1677,18 @@ public final class FleetMcp {
|
||||
* an explicit {@code true}
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
|
||||
outage, leadSeats, leads, selfTerm, coordination, callerIsPrimary);
|
||||
}
|
||||
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
LoopHealthSource loopHealth,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
try {
|
||||
@@ -1656,6 +1712,9 @@ public final class FleetMcp {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
result.put("leads", leadRows); result.put("members", out);
|
||||
result.put("healthCoverage", healthCoverage.value().get());
|
||||
result.put("loopHealth", Map.of(
|
||||
"statusPoller", loopHealth.statusPoller().get().name(),
|
||||
"sessionReaper", loopHealth.sessionReaper().get().name()));
|
||||
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
|
||||
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
|
||||
// key is absent rather than present-and-empty.
|
||||
@@ -2135,9 +2194,10 @@ public final class FleetMcp {
|
||||
+ "cannot reliably re-identify: some backends (e.g. opencode) resolve it from the "
|
||||
+ "member's working directory, which only uniquely identifies a member when it "
|
||||
+ "was spawned into its own fleetd-provisioned worktree (worktree:true/<slug>); a "
|
||||
+ "member spawned without one shares its directory with others and never reports "
|
||||
+ "an id, however long it runs (fleetd #249). An empty 'members' "
|
||||
+ "means no members are spawned; it says nothing about peers. When capacity "
|
||||
+ "member spawned without one shares its directory with others and never reports "
|
||||
+ "an id, however long it runs (fleetd #249). An empty 'members' "
|
||||
+ "means no members are spawned; it says nothing about peers. 'loopHealth' reports "
|
||||
+ "the RUNNING, STALLED, or STOPPED state of statusPoller and sessionReaper. When capacity "
|
||||
+ "facts are configured, a 'capacity' row per profile reports 'free' — the "
|
||||
+ "slots a fresh fleet_spawn on that profile will actually be granted right "
|
||||
+ "now (max(0, maxLoad - live)), the same check the spawn gate itself runs. A "
|
||||
|
||||
@@ -93,19 +93,30 @@ public final class MessageService {
|
||||
/** Timed out after the message was delivered — the worker is still working. */
|
||||
TIMED_OUT_WORKING,
|
||||
/**
|
||||
* Timed out with no confirmed delivery. Despite the name, this does not mean the message
|
||||
* is sitting in a queue — on every route it will not arrive later. {@link #send} reaches
|
||||
* this outcome through {@link Injector#cancel}, whose result tells the routes apart:
|
||||
* {@link Injector.Cancellation#CANCELLED} means the message was still queued and this call
|
||||
* removed it, so the target saw nothing; {@link Injector.Cancellation#NOT_DELIVERED} means
|
||||
* an earlier attempt already decided the message's fate — the injector's call to the
|
||||
* target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) threw, or the queue
|
||||
* was cleared because the target never became ready or was abandoned — and {@code cancel}
|
||||
* is only reporting that pre-existing state. On the failed-attempt route the terminal may
|
||||
* already hold a partial paste from before the call threw, so only the {@code CANCELLED}
|
||||
* case establishes that the target saw nothing.
|
||||
* Timed out with no confirmed delivery, and the target saw nothing — the message will not
|
||||
* arrive later, so a caller may resend. {@link #send} reaches this outcome through {@link
|
||||
* Injector#cancel} reporting one of two routes: {@link Injector.Cancellation#CANCELLED}
|
||||
* means the message was still queued and this call removed it; {@link
|
||||
* Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the queue was cleared
|
||||
* because the target never became ready or was abandoned, or the injector's call to the
|
||||
* target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) failed with a herdr
|
||||
* error that this codebase already treats as a confirmed absence. A third route,
|
||||
* {@link Injector.Cancellation#ATTEMPTED}, used to be folded into this same outcome
|
||||
* (fleetd #571) — it no longer is; see {@link #TIMED_OUT_UNCONFIRMED}.
|
||||
*/
|
||||
TIMED_OUT_QUEUED,
|
||||
/**
|
||||
* Timed out with delivery unknown. {@link #send} reaches this outcome when {@link
|
||||
* Injector#cancel} reports {@link Injector.Cancellation#ATTEMPTED} (fleetd #551): the call
|
||||
* to the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) was made, but
|
||||
* this caller never observed whether it reached the pane. {@code agent.prompt} pastes
|
||||
* <em>and submits</em> in one call, so the target may already hold a complete, submitted
|
||||
* turn and be working on it right now — the same reality as {@link #TIMED_OUT_WORKING},
|
||||
* just not confirmed. The message may or may not have arrived. Treat this as neither a
|
||||
* confirmed delivery nor a confirmed absence: a caller that resends on this outcome risks a
|
||||
* double delivery — the same brief typed into the pane twice (fleetd #571).
|
||||
*/
|
||||
TIMED_OUT_UNCONFIRMED,
|
||||
/** Another send to this session was in flight for the whole window. */
|
||||
BUSY,
|
||||
/**
|
||||
@@ -311,18 +322,21 @@ public final class MessageService {
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Boolean> strandedReplies = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Targets whose last send timed out with {@link Outcome#TIMED_OUT_QUEUED} (CB-640) — {@link
|
||||
* #send} called {@link Injector#cancel} and got back something other than {@code DELIVERED}.
|
||||
* That covers more than one history: the message may have still been queued and {@code cancel}
|
||||
* removed it right there, or an earlier attempt may already have failed (the call to the
|
||||
* target's terminal threw) or been abandoned (the target never became ready, or was torn
|
||||
* down). What holds on every route: the message will not arrive later, and it is not sitting
|
||||
* in any queue. What does NOT hold on every route: that the target saw nothing — a failed
|
||||
* delivery attempt can leave a partial paste behind. Set where {@link #send} already computes
|
||||
* {@code wasDelivered} for that outcome; no queue is kept here, only the fact that the send
|
||||
* ended with no confirmed delivery. Cleared the same way as {@link #strandedReplies}: the next
|
||||
* accepted delivery for the target ({@link #send} opening a fresh waiter) or a teardown
|
||||
* ({@link #abandon}).
|
||||
* Targets whose last send timed out with no confirmed delivery (CB-640) — {@link #send} called
|
||||
* {@link Injector#cancel} and got back something other than {@code DELIVERED}. That covers
|
||||
* three histories, not one: {@link Injector.Cancellation#CANCELLED} — the message was still
|
||||
* queued and {@code cancel} removed it right there; {@link Injector.Cancellation#NOT_DELIVERED}
|
||||
* — nothing was ever sent, because the target never became ready, was torn down, or the call to
|
||||
* its terminal failed with a herdr error this codebase already treats as a confirmed absence; or
|
||||
* {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) — the call to the target's terminal was
|
||||
* made and its outcome is unknown, so the target may already hold a complete, submitted turn.
|
||||
* Only the first two mean the message will not arrive later and the target saw nothing; on the
|
||||
* third it may already have arrived in full — and the caller sees a different outcome for it
|
||||
* ({@link Outcome#TIMED_OUT_UNCONFIRMED}, fleetd #571) than for the first two ({@link
|
||||
* Outcome#TIMED_OUT_QUEUED}). Set where {@link #send} already computes {@code wasDelivered} for
|
||||
* that outcome; no queue is kept here, only the fact that the send ended with no confirmed
|
||||
* delivery. Cleared the same way as {@link #strandedReplies}: the next accepted delivery for the
|
||||
* target ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}).
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
|
||||
private final AtomicLong ticketSeq = new AtomicLong();
|
||||
@@ -409,17 +423,22 @@ public final class MessageService {
|
||||
|
||||
/**
|
||||
* Read-only delegation fact for fleet views (CB-640): {@code target}'s last send timed out
|
||||
* with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} (see the
|
||||
* {@code TimeoutException} branch of {@link #send}). Despite the method's name, this is not
|
||||
* proof that a message is sitting in a queue: {@link Injector#cancel} reports this outcome
|
||||
* either because the message was still queued and got removed right there, or because an
|
||||
* earlier attempt already failed (the call to the target's terminal threw) or was abandoned
|
||||
* (the target never became ready, or was torn down). Either way the message will not arrive
|
||||
* later. It does NOT follow that the target saw nothing — on the failed-attempt route the
|
||||
* terminal may already hold a partial paste. Distinct from {@link Outcome#TIMED_OUT_WORKING},
|
||||
* where delivery already happened and only the reply is outstanding. Cleared the next time
|
||||
* this target's delivery is accepted or the target is abandoned — see
|
||||
* {@link #queuedDeliveries}.
|
||||
* with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} or {@link
|
||||
* Outcome#TIMED_OUT_UNCONFIRMED} (fleetd #571; see the {@code TimeoutException} branch of
|
||||
* {@link #send}). Despite the method's name, this is not proof that a message is sitting in a
|
||||
* queue: {@link Injector#cancel} reports this outcome through three routes. {@link
|
||||
* Injector.Cancellation#CANCELLED} means the message was still queued and got removed right
|
||||
* there. {@link Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the target
|
||||
* never became ready, was torn down, or the call to its terminal failed with a herdr error this
|
||||
* codebase already treats as a confirmed absence. Only these two routes mean the message will
|
||||
* not arrive later, and both report {@code TIMED_OUT_QUEUED}. {@link
|
||||
* Injector.Cancellation#ATTEMPTED} (fleetd #551) means the call to the target's terminal was
|
||||
* made and its outcome is unknown: {@code agent.prompt} pastes <em>and submits</em> in one
|
||||
* call, so on this route the target may already hold a complete, submitted turn and be
|
||||
* working on it right now — it does NOT follow that the target saw nothing, and this route
|
||||
* reports {@code TIMED_OUT_UNCONFIRMED} instead. Distinct from {@link Outcome#TIMED_OUT_WORKING},
|
||||
* where delivery already happened and only the reply is outstanding. Cleared the next time this
|
||||
* target's delivery is accepted or the target is abandoned — see {@link #queuedDeliveries}.
|
||||
*/
|
||||
public boolean hasQueuedDelivery(String target) {
|
||||
return target != null && queuedDeliveries.containsKey(target);
|
||||
@@ -634,7 +653,7 @@ public final class MessageService {
|
||||
return switch (o) {
|
||||
case REPLIED -> "replied";
|
||||
case COMPLETED_UNREPLIED -> "completion_fallback";
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
|
||||
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
|
||||
@@ -963,6 +982,7 @@ public final class MessageService {
|
||||
} catch (TimeoutException e) {
|
||||
boolean wasDelivered = delivery.completion().isDone()
|
||||
&& !delivery.completion().isCompletedExceptionally();
|
||||
Injector.Cancellation cancellation = null;
|
||||
if (!wasDelivered) {
|
||||
if (timeoutCancellationRaceHookForTest != null) {
|
||||
// Test-only (fleetd #345): see the field's own javadoc.
|
||||
@@ -970,19 +990,27 @@ public final class MessageService {
|
||||
}
|
||||
// The target monitor makes cancellation atomic with onStatus picking this
|
||||
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
|
||||
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
|
||||
cancellation = injector.cancel(delivery);
|
||||
wasDelivered = cancellation == Injector.Cancellation.DELIVERED;
|
||||
}
|
||||
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
|
||||
if (!wasDelivered) {
|
||||
// CB-640: record that delivery did not happen for fleet health (see
|
||||
// queuedDeliveries). Whatever injector.cancel() reported above — this call
|
||||
// removed a still-queued Pending, or an earlier attempt already failed or
|
||||
// was abandoned — the send ends with no confirmed delivery and will not
|
||||
// arrive later.
|
||||
Outcome outcome;
|
||||
if (wasDelivered) {
|
||||
outcome = Outcome.TIMED_OUT_WORKING;
|
||||
} else if (cancellation == Injector.Cancellation.ATTEMPTED) {
|
||||
// fleetd #571: the call to the target's terminal was made and its outcome is
|
||||
// unknown — the message may already have arrived in full, so this must not
|
||||
// be reported as TIMED_OUT_QUEUED, which promises it never will.
|
||||
outcome = Outcome.TIMED_OUT_UNCONFIRMED;
|
||||
} else {
|
||||
// CB-640: record that delivery is not confirmed, for fleet health (see
|
||||
// queuedDeliveries). cancellation is CANCELLED (this call removed a
|
||||
// still-queued Pending) or NOT_DELIVERED (an earlier attempt already failed
|
||||
// with a confirmed absence) — both mean the target saw nothing.
|
||||
queuedDeliveries.put(target, Boolean.TRUE);
|
||||
outcome = Outcome.TIMED_OUT_QUEUED;
|
||||
}
|
||||
return recorded(new Reply(
|
||||
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
|
||||
return recorded(new Reply(outcome, null));
|
||||
} catch (ExecutionException e) {
|
||||
Throwable cause = e.getCause();
|
||||
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
||||
@@ -1122,21 +1150,33 @@ public final class MessageService {
|
||||
}
|
||||
try {
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
|
||||
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
|
||||
// resumed turn can re-associate the async ticket with its new turnId via
|
||||
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
|
||||
// markAsyncQuestion silently returns null.
|
||||
Task task = asyncTasksByTurn.get(turnId);
|
||||
if (task != null) {
|
||||
asyncTasksByWaiter.put(reply, task);
|
||||
}
|
||||
if (!rendezvous.answerAsk(turnId, content)) {
|
||||
asyncTasksByWaiter.remove(reply);
|
||||
rendezvous.close(workerSession, reply);
|
||||
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
|
||||
}
|
||||
clearAsyncQuestion(turnId, false);
|
||||
// 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
|
||||
// try up to wrap the registration closes the gap structurally: every exit from here on,
|
||||
// STALE_TURN included, now runs through the one finally below exactly once, and the
|
||||
// duplicated pair is gone. This did not fix a live leak — see the ticket: neither
|
||||
// rendezvous.answerAsk nor clearAsyncQuestion(turnId, false) can throw, so nothing ever
|
||||
// actually left through the old gap uncovered — but #572 found this exact drift (one
|
||||
// finally asserted, an identical sibling not) on this same file, and the hand-rolled copy
|
||||
// was the wrong shape to keep regardless of whether it was ever exercised.
|
||||
try {
|
||||
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
|
||||
// resumed turn can re-associate the async ticket with its new turnId via
|
||||
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
|
||||
// markAsyncQuestion silently returns null.
|
||||
Task task = asyncTasksByTurn.get(turnId);
|
||||
if (task != null) {
|
||||
asyncTasksByWaiter.put(reply, task);
|
||||
}
|
||||
if (answerAskLapseRaceHookForTest != null) {
|
||||
// Test-only (fleetd #575): see the field's own javadoc.
|
||||
answerAskLapseRaceHookForTest.run();
|
||||
}
|
||||
if (!rendezvous.answerAsk(turnId, content)) {
|
||||
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
|
||||
}
|
||||
clearAsyncQuestion(turnId, false);
|
||||
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
||||
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
|
||||
// #282: this waiter can resolve with a FRESH question rather than a terminal reply —
|
||||
@@ -1636,6 +1676,30 @@ public final class MessageService {
|
||||
this.askTimeoutRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #575 — invoked from {@link #answer}, right after this
|
||||
* call's own {@code Task} registration and right before its {@code rendezvous.answerAsk(turnId,
|
||||
* content)} call. A test installs this to complete the SAME turnId's ask directly via {@link
|
||||
* Rendezvous#answerAsk} from inside that exact window, deterministically reproducing what a
|
||||
* second, concurrent {@code answer()} call racing to unblock the same ask can otherwise only
|
||||
* win by timing luck: this call's own {@code askSession(turnId)} lookup at the top already saw
|
||||
* the ask as open, but by the time it reaches {@code rendezvous.answerAsk} here, the other call
|
||||
* already completed it (or the worker's own {@code ask()} teardown already closed it) — so this
|
||||
* call must see {@code false} and return {@link Outcome#STALE_TURN}, exactly the "lapsed between
|
||||
* the lookup and the unblock" case named at that call site. Proves the #575 fix (widening this
|
||||
* method's try so a single finally covers this exit) does not change that outcome and still
|
||||
* cleans this call's own {@code reply} up exactly once.
|
||||
*/
|
||||
private volatile Runnable answerAskLapseRaceHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #575): install {@link #answerAskLapseRaceHookForTest}. Package-private so the
|
||||
* test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setAnswerAskLapseRaceHookForTest(Runnable hook) {
|
||||
this.answerAskLapseRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
|
||||
private boolean hasAsyncQuestion(String target) {
|
||||
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
|
||||
|
||||
@@ -94,6 +94,7 @@ public final class FleetApp {
|
||||
// .none() (the honest "feature not wired" view) for every constructor that does not pass one.
|
||||
private final FleetMcp.QuarantineSource quarantine;
|
||||
private final FleetMcp.OutageSource outage;
|
||||
private final FleetMcp.LoopHealthSource loopHealth;
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
|
||||
/**
|
||||
@@ -155,7 +156,8 @@ public final class FleetApp {
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials) {
|
||||
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics,
|
||||
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
|
||||
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LoopHealthSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -169,8 +171,18 @@ public final class FleetApp {
|
||||
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
|
||||
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
|
||||
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
|
||||
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, deliverable,
|
||||
memberCredentials, quarantine, outage, FleetMcp.LoopHealthSource.none());
|
||||
}
|
||||
|
||||
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
|
||||
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage,
|
||||
FleetMcp.LoopHealthSource loopHealth) {
|
||||
this.herdr = herdr;
|
||||
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
|
||||
this.workers = workers;
|
||||
@@ -183,6 +195,7 @@ public final class FleetApp {
|
||||
this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent;
|
||||
this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none();
|
||||
this.outage = outage != null ? outage : FleetMcp.OutageSource.none();
|
||||
this.loopHealth = loopHealth != null ? loopHealth : FleetMcp.LoopHealthSource.none();
|
||||
}
|
||||
|
||||
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
|
||||
@@ -291,15 +304,19 @@ public final class FleetApp {
|
||||
* spotted by comparing two numbers by eye.
|
||||
*/
|
||||
private void healthz(Context ctx) {
|
||||
HealthzResponse response = healthzResponse(herdr, memberHerdr, loopHealth);
|
||||
ctx.status(response.status()).json(response.body());
|
||||
}
|
||||
|
||||
record HealthzResponse(int status, Map<String, Object> body) { }
|
||||
|
||||
static HealthzResponse healthzResponse(HerdrClient herdr, HerdrClient memberHerdr,
|
||||
FleetMcp.LoopHealthSource loopHealth) {
|
||||
JsonNode pong;
|
||||
try {
|
||||
pong = herdr.call("ping");
|
||||
} catch (HerdrException e) {
|
||||
ctx.status(503).json(Map.of(
|
||||
"status", "degraded",
|
||||
"herdr", "unreachable",
|
||||
"detail", e.getMessage()));
|
||||
return;
|
||||
return degradedResponse("unreachable", e.getMessage(), loopHealth);
|
||||
}
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
body.put("status", "ok");
|
||||
@@ -311,11 +328,11 @@ public final class FleetApp {
|
||||
try {
|
||||
memberPong = memberHerdr.call("ping");
|
||||
} catch (HerdrException e) {
|
||||
ctx.status(503).json(Map.of(
|
||||
return new HealthzResponse(503, Map.of(
|
||||
"status", "degraded",
|
||||
"herdr", "member unreachable",
|
||||
"detail", e.getMessage()));
|
||||
return;
|
||||
"detail", e.getMessage(),
|
||||
"loopHealth", loopHealthView(loopHealth)));
|
||||
}
|
||||
int leadProtocol = pong.path("protocol").asInt();
|
||||
int memberProtocol = memberPong.path("protocol").asInt();
|
||||
@@ -326,7 +343,20 @@ public final class FleetApp {
|
||||
body.put("protocolMismatch", true);
|
||||
}
|
||||
}
|
||||
ctx.status(200).json(body);
|
||||
body.put("loopHealth", loopHealthView(loopHealth));
|
||||
return new HealthzResponse(200, body);
|
||||
}
|
||||
|
||||
private static Map<String, String> loopHealthView(FleetMcp.LoopHealthSource loopHealth) {
|
||||
return Map.of("statusPoller", loopHealth.statusPoller().get().name(),
|
||||
"sessionReaper", loopHealth.sessionReaper().get().name());
|
||||
}
|
||||
|
||||
private static HealthzResponse degradedResponse(String herdr, String detail,
|
||||
FleetMcp.LoopHealthSource loopHealth) {
|
||||
return new HealthzResponse(503, Map.of(
|
||||
"status", "degraded", "herdr", herdr, "detail", detail,
|
||||
"loopHealth", loopHealthView(loopHealth)));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -640,15 +670,31 @@ public final class FleetApp {
|
||||
}
|
||||
default -> ctx.status(202).json(Map.of(
|
||||
"sessionId", id,
|
||||
// fleetd #571 (ticket comment 17126): no `default` here on purpose. This switch
|
||||
// is an expression, so the compiler already demands every Outcome constant have
|
||||
// an arm — adding an 11th constant to Outcome is a compile error here, not a
|
||||
// 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.
|
||||
"status", switch (reply.outcome()) {
|
||||
case TIMED_OUT_WORKING -> "working";
|
||||
case TIMED_OUT_QUEUED -> "queued";
|
||||
// Delivery here is unknown, not merely still queued — see
|
||||
// Outcome#TIMED_OUT_UNCONFIRMED's own javadoc.
|
||||
case TIMED_OUT_UNCONFIRMED -> "unconfirmed";
|
||||
case BUSY -> "busy";
|
||||
case WORKER_FAILED -> "failed";
|
||||
case BACKEND_EXHAUSTED -> "backend_exhausted";
|
||||
default -> "done"; // unreachable (terminal outcomes handled above)
|
||||
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable
|
||||
},
|
||||
"detail", (reply.outcome() == MessageService.Outcome.WORKER_FAILED
|
||||
"detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED
|
||||
? "no reply within " + timeout + "ms; delivery is unconfirmed — the "
|
||||
+ "message may already have reached the worker, so a resend "
|
||||
+ "risks sending it twice; poll status first"
|
||||
: (reply.outcome() == MessageService.Outcome.WORKER_FAILED
|
||||
|| reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED)
|
||||
&& reply.text() != null
|
||||
? reply.text()
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.function.BiConsumer;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :591}. {@code Fleetd.main} wires {@link
|
||||
* dev.ltms.fleet.health.FleetHealthMonitor}'s {@code failTarget} callback with {@code
|
||||
* messages::abandon} — before this ticket that was an inline argument to {@code new
|
||||
* FleetHealthMonitor(...)}. Measured: replacing it with a no-op {@code BiConsumer} at the call
|
||||
* site compiles with 0 errors and leaves the full suite green, because nothing else in the tree
|
||||
* ever drives that specific constructor argument. In production it means a member the monitor
|
||||
* classifies {@code GONE}/{@code NEVER_READY} never has its pending ticket failed — the caller
|
||||
* keeps reporting {@code PENDING} for the full 30-minute async timeout instead of the immediate,
|
||||
* accurate failure CB-580 exists to give it.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#healthFailTarget} directly — never {@code FleetHealthMonitor}
|
||||
* or {@code main} — against a real {@link MessageService}, using the same {@code sendAsync} +
|
||||
* {@code poll} observable {@link MessageServiceTest} already relies on to pin {@code
|
||||
* MessageService.abandon} itself.
|
||||
*/
|
||||
class FleetdHealthFailTargetWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.healthFailTarget delegates to the real MessageService.abandon, not a no-op")
|
||||
void healthFailTargetDelegatesToMessagesAbandon() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("$ prompt");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous);
|
||||
|
||||
BiConsumer<String, String> failTarget = Fleetd.healthFailTarget(messages);
|
||||
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting(rendezvous);
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
|
||||
|
||||
failTarget.accept(T, "member unreachable (health monitor)");
|
||||
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() != MessageService.Phase.PENDING) {
|
||||
break;
|
||||
}
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(),
|
||||
"Fleetd.healthFailTarget(messages) must return messages::abandon — replacing it "
|
||||
+ "with a no-op BiConsumer at the Fleetd.healthFailTarget call site means "
|
||||
+ "this ticket is never failed and keeps polling as PENDING");
|
||||
assertTrue(view.detail() != null && view.detail().contains("member unreachable"),
|
||||
"the failure reason passed to failTarget.accept must reach MessageService.abandon "
|
||||
+ "and end up in the ticket's detail");
|
||||
}
|
||||
|
||||
private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.net.ServerSocket;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :518}. {@code Fleetd.main} passes {@code LeadMailbox::open} as
|
||||
* the {@link Fleetd.LeadMailboxOpener} argument to {@code openLeadMailbox(...)} — before this
|
||||
* ticket that method reference was inline at the call site. Measured: replacing it with the inert
|
||||
* {@code (uri, selfCoordId, prefetch) -> null} compiles with 0 errors and leaves the full suite
|
||||
* green — {@code FleetdLeadMailboxSelectionTest} drives {@code openLeadMailbox} with its own
|
||||
* injected opener and never observes what {@code main} itself actually passes.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#leadMailboxOpener} directly and proves it is the real,
|
||||
* network-attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must
|
||||
* throw, exactly mirroring {@link FleetdReplyInboxOpenerWiringTest} for the reply-inbox opener.
|
||||
* The inert form never attempts a connection and returns {@code null} without throwing, so it
|
||||
* fails this assertion silently.
|
||||
*/
|
||||
class FleetdLeadMailboxOpenerWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.leadMailboxOpener is the real LeadMailbox::open, not a stub that never connects")
|
||||
void leadMailboxOpenerAttemptsARealConnection() throws Exception {
|
||||
int closedPort;
|
||||
try (ServerSocket socket = new ServerSocket(0)) {
|
||||
closedPort = socket.getLocalPort();
|
||||
} // released immediately — connecting to it now is a guaranteed refusal, not a fluke
|
||||
|
||||
Fleetd.LeadMailboxOpener opener = Fleetd.leadMailboxOpener();
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/coord", "coord-1", 50),
|
||||
"Fleetd.leadMailboxOpener() must be LeadMailbox::open — a real network attempt "
|
||||
+ "against a genuinely unreachable broker must throw. The inert form "
|
||||
+ "(uri, selfCoordId, prefetch) -> null never attempts a connection and "
|
||||
+ "returns null instead of throwing, so it would fail this assertion "
|
||||
+ "silently.");
|
||||
assertTrue(thrown.getMessage().contains("cannot connect to AMQP coordination broker"),
|
||||
"must be LeadMailbox.open's own real failure message, not a different exception "
|
||||
+ "shape standing in for it");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,174 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.LoopWatchdog;
|
||||
import dev.ltms.fleet.inject.StatusPoller;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.SessionReaper;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #562 follow-up (issue comment "HOLD on PR #579"): {@code Fleetd.main}'s {@code loopHealth}
|
||||
* local used to be a bare {@code new FleetMcp.LoopHealthSource(poller::health, ...)} built inline,
|
||||
* with nothing a test could call directly. Measured on that shape: replacing {@code
|
||||
* poller::health} with a constant {@code () -> LoopWatchdog.State.RUNNING} at the call site
|
||||
* compiled with 0 errors and left all 1771 existing tests green — the daemon could be changed to
|
||||
* always report the {@link StatusPoller} as {@code RUNNING}, so the watchdog could never fire and
|
||||
* a stalled poller would be invisible, while every test stayed green. That is exactly the false
|
||||
* negative this ticket exists to prevent.
|
||||
*
|
||||
* <p>The five tests PR #579 added ({@code FleetMcpTest}, {@code FleetAppTest}) all build their own
|
||||
* {@link FleetMcp.LoopHealthSource} directly with fixed lambdas — they prove the seam ({@code
|
||||
* LoopHealthSource} reports what it is given) and nothing about what {@code Fleetd.main} actually
|
||||
* gives it. This is the same hand-built-vs-config-wired shape as fleetd #561/#248/#426.
|
||||
*
|
||||
* <p>The fix extracts the inline {@code new} into {@link Fleetd#loopHealthSource}, a package-private
|
||||
* factory in the same style as {@link Fleetd#capacitySource} and {@link Fleetd#healthCoverageSource}
|
||||
* — which is exactly what makes it directly callable here. This test calls that factory with real
|
||||
* {@link StatusPoller}/{@link SessionReaper} instances (never started, so no herdr or git I/O
|
||||
* happens) and pins each half separately, plus the {@code reaper == null} branch: one invariant
|
||||
* wired at three places needs three assertions, not one combined check whose non-zero total could
|
||||
* hide a gap at any single place.
|
||||
*/
|
||||
class FleetdLoopHealthSourceWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("the statusPoller half reports the real poller's health, not a hardcoded state")
|
||||
void statusPollerHalfReflectsThePollersRealHealth() {
|
||||
// Stopped without ever being started — stop() still marks the watchdog STOPPED. A poller
|
||||
// that has never reported RUNNING is the discriminating case: if Fleetd.loopHealthSource
|
||||
// ever hardcoded RUNNING (the exact mutation this test exists to catch), this would fail.
|
||||
StatusPoller stoppedPoller = freshPoller();
|
||||
stoppedPoller.stop();
|
||||
SessionReaper unusedReaper = freshReaper(); // present only to satisfy the signature
|
||||
|
||||
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(stoppedPoller, unusedReaper);
|
||||
|
||||
assertEquals(LoopWatchdog.State.STOPPED, source.statusPoller().get(),
|
||||
"the statusPoller supplier must delegate to the real poller's health() — "
|
||||
+ "replacing poller::health with a constant () -> RUNNING at the "
|
||||
+ "Fleetd.loopHealthSource call site must fail this assertion");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the sessionReaper half reports the real reaper's health, not a hardcoded state")
|
||||
void sessionReaperHalfReflectsTheReapersRealHealth() {
|
||||
StatusPoller unusedPoller = freshPoller(); // present only to satisfy the signature
|
||||
SessionReaper stoppedReaper = freshReaper();
|
||||
stoppedReaper.stop();
|
||||
|
||||
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(unusedPoller, stoppedReaper);
|
||||
|
||||
assertEquals(LoopWatchdog.State.STOPPED, source.sessionReaper().get(),
|
||||
"the sessionReaper supplier must delegate to the real reaper's health() — "
|
||||
+ "replacing reaper.health() with a constant at the "
|
||||
+ "Fleetd.loopHealthSource call site must fail this assertion");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a null reaper (idle ttl not configured) still reports STOPPED, not a crash")
|
||||
void nullReaperStillReportsStopped() {
|
||||
// SessionReaper is only constructed when lifecycle.idleTtlSeconds is configured (see the
|
||||
// `reaper` local in Fleetd.main) — a real deployment routinely passes null here. That null
|
||||
// check is real behaviour, not a simplification to delete: it must keep reporting STOPPED
|
||||
// rather than throwing a NullPointerException on the first fleet_list/healthz call.
|
||||
StatusPoller runningPoller = freshPoller();
|
||||
|
||||
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(runningPoller, null);
|
||||
|
||||
assertEquals(LoopWatchdog.State.STOPPED, source.sessionReaper().get(),
|
||||
"reaper == null must still report STOPPED, exactly like an intentionally-stopped "
|
||||
+ "reaper would — do not delete this null check to simplify the wiring");
|
||||
}
|
||||
|
||||
/** Never started, so no herdr call is ever made; freshly constructed reports RUNNING. */
|
||||
private static StatusPoller freshPoller() {
|
||||
AgentControl agents = new AgentControl(new FakeHerdr());
|
||||
return new StatusPoller(agents, new Injector(agents), 1000);
|
||||
}
|
||||
|
||||
/** Never started, so no git/session I/O is ever made; freshly constructed reports RUNNING. */
|
||||
private static SessionReaper freshReaper() {
|
||||
return new SessionReaper(new SessionManager(new NeverSpawnsLauncher()), 60, 1000);
|
||||
}
|
||||
|
||||
/**
|
||||
* Same minimal shape as {@code FleetdBackendErrorSinkTest.NeverSpawnsLauncher} — every method
|
||||
* throws or returns an empty/no-op value, since a {@link SessionReaper} that is only ever
|
||||
* constructed and then stopped (never started) never calls any of them.
|
||||
*/
|
||||
private static final class NeverSpawnsLauncher implements PeerLauncher {
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return Set.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilitiesFor(String profileName) {
|
||||
return Set.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> parityOverlay(String profileName) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<?> list() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean clearContext(String id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :611-631}. {@code Fleetd.main} wires {@code
|
||||
* sessions.onRelease(...)} with a lambda that calls three collaborators — {@code
|
||||
* messages.abandon}, {@code replyInbox.release}, and {@code primaryRegistry.forgetDelegation} —
|
||||
* before this ticket built inline inside {@code main}. Measured: replacing the whole lambda body
|
||||
* with {@code detail -> { }} compiles with 0 errors and leaves the full suite green, because each
|
||||
* collaborator is separately tested in isolation ({@code MessageServiceTest}, {@code
|
||||
* PrimaryRegistryTest}) but nothing before this ticket drove the lambda that calls all three from
|
||||
* {@code main}.
|
||||
*
|
||||
* <p>{@code MessageService.abandon}'s own javadoc documents the consequence: without this,
|
||||
* tearing a worker down leaves its rendezvous waiter open, so a blocking {@code fleet_send} keeps
|
||||
* blocking and an async one reports {@code PENDING} for a hardcoded thirty minutes on every {@code
|
||||
* fleet_stop} and every idle-reap.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#releaseCleanup} directly — never {@code SessionManager} or
|
||||
* {@code main} — against real {@link MessageService}, {@link InMemoryReplyInbox}, and {@link
|
||||
* PrimaryRegistry} instances, and asserts each collaborator's own observable effect: the pending
|
||||
* async ticket transitions to {@code FAILED} (abandon), the inbox no longer owns the target's
|
||||
* queue (release), and the recorded delegation is forgotten (forgetDelegation).
|
||||
*/
|
||||
class FleetdReleaseCleanupWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.releaseCleanup reaches messages.abandon, replyInbox.release, and primaryRegistry.forgetDelegation")
|
||||
void releaseCleanupReachesAllThreeCollaborators() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("$ prompt");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous);
|
||||
InMemoryReplyInbox replyInbox = new InMemoryReplyInbox();
|
||||
PrimaryRegistry primaryRegistry = new PrimaryRegistry(null);
|
||||
|
||||
// Set up the "before" state each collaborator's own effect is measured against.
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting(rendezvous);
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
|
||||
"sanity: the async ticket is pending before cleanup runs");
|
||||
|
||||
replyInbox.own(T);
|
||||
replyInbox.publish(T, "msg-1", "hello");
|
||||
assertEquals(1, replyInbox.peek(T).size(),
|
||||
"sanity: the inbox owns T and holds one message before cleanup runs");
|
||||
|
||||
primaryRegistry.recordDelegation(T, "lead-1");
|
||||
assertEquals("lead-1", primaryRegistry.nudgeTargetFor(T).orElse(null),
|
||||
"sanity: the delegation is recorded before cleanup runs");
|
||||
|
||||
Consumer<SessionManager.ReleaseDetail> cleanup =
|
||||
Fleetd.releaseCleanup(messages, replyInbox, primaryRegistry);
|
||||
cleanup.accept(new SessionManager.ReleaseDetail(T, null, null, null, null));
|
||||
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() != MessageService.Phase.PENDING) {
|
||||
break;
|
||||
}
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(),
|
||||
"releaseCleanup must call messages.abandon(...) — an inert detail -> { } lambda "
|
||||
+ "leaves this ticket PENDING forever");
|
||||
assertTrue(view.detail() != null && view.detail().contains("released"),
|
||||
"the abandon reason must say the worker session was released");
|
||||
|
||||
assertTrue(replyInbox.peek(T).isEmpty(),
|
||||
"releaseCleanup must call replyInbox.release(...) — an inert lambda leaves the "
|
||||
+ "inbox still owning T with its message");
|
||||
|
||||
assertTrue(primaryRegistry.nudgeTargetFor(T).isEmpty(),
|
||||
"releaseCleanup must call primaryRegistry.forgetDelegation(...) — an inert lambda "
|
||||
+ "leaves the stale delegation in place");
|
||||
}
|
||||
|
||||
private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.net.ServerSocket;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :512}. {@code Fleetd.main} passes {@code AmqpReplyInbox::open}
|
||||
* as the {@link Fleetd.AmqpOpener} argument to {@code selectReplyInbox(...)} — before this ticket
|
||||
* that method reference was inline at the call site. Measured: replacing it with the inert {@code
|
||||
* (uri, prefetch) -> new InMemoryReplyInbox()} compiles with 0 errors and leaves the full suite
|
||||
* green — {@code FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own
|
||||
* injected opener (including one test that passes the real {@code AmqpReplyInbox::open}
|
||||
* explicitly) and never observes what {@code main} itself actually passes.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#replyInboxOpener} directly and proves it is the real, network-
|
||||
* attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must throw —
|
||||
* the same shape {@code FleetdReplyInboxSelectionTest.aRealUnreachableBrokerFallsBackViaTheRealOpener}
|
||||
* already relies on for {@code AmqpReplyInbox.open} itself. The inert form never attempts a
|
||||
* connection and never throws, so it fails this assertion silently (by returning normally).
|
||||
*/
|
||||
class FleetdReplyInboxOpenerWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.replyInboxOpener is the real AmqpReplyInbox::open, not a stub that never connects")
|
||||
void replyInboxOpenerAttemptsARealConnection() throws Exception {
|
||||
int closedPort;
|
||||
try (ServerSocket socket = new ServerSocket(0)) {
|
||||
closedPort = socket.getLocalPort();
|
||||
} // released immediately — connecting to it now is a guaranteed refusal, not a fluke
|
||||
|
||||
Fleetd.AmqpOpener opener = Fleetd.replyInboxOpener();
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/vh", 50),
|
||||
"Fleetd.replyInboxOpener() must be AmqpReplyInbox::open — a real network attempt "
|
||||
+ "against a genuinely unreachable broker must throw. The inert form "
|
||||
+ "(uri, prefetch) -> new InMemoryReplyInbox() never attempts a connection "
|
||||
+ "and never throws, so it would return normally here and fail this "
|
||||
+ "assertion silently.");
|
||||
assertTrue(thrown.getMessage().contains("cannot connect to AMQP broker"),
|
||||
"must be AmqpReplyInbox.open's own real failure message, not a different exception "
|
||||
+ "shape standing in for it");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,286 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.inject.TurnListener;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #561: {@code Fleetd.turnListener} composes the completion resolver and the session
|
||||
* manager into one {@link TurnListener}. Each of the four callbacks below used to be two bare,
|
||||
* unguarded statements (completion first, then sessions) — a throw from the session half used to
|
||||
* skip nothing <em>after</em> it (there was nothing after it), but nothing enforced that the
|
||||
* completion half had to come first either, beyond call order in the source. {@code onDelivered}
|
||||
* had the identical shape and was fixed by #556 (moving its registration off the fan-out
|
||||
* entirely); these four callbacks cannot be fixed that way, so the fan-out itself is hardened
|
||||
* instead (see {@link Fleetd#turnListener} and {@code Fleetd.bothMustRun}/{@code
|
||||
* bothMustRunKeepingSecondResult}).
|
||||
*
|
||||
* <p>The invariant under test: <b>a throwing session-listener must not prevent the completion
|
||||
* resolver from being told the turn ended.</b> Each test below builds the exact production
|
||||
* composition ({@link Fleetd#turnListener}) from a real {@link CompletionResolver} and a fake
|
||||
* {@code sessions} half that throws, then asserts the completion half's effect (the captured
|
||||
* {@link Rendezvous} waiter resolving) happened anyway — never by inspecting call order directly.
|
||||
*/
|
||||
class FleetdTurnListenerCompositionTest {
|
||||
|
||||
/** Records which callbacks ran and can be told to throw from a chosen one. */
|
||||
private static final class RecordingSessions implements TurnListener {
|
||||
final Set<String> called = new LinkedHashSet<>();
|
||||
private final Set<String> throwing;
|
||||
|
||||
RecordingSessions(String... throwingMethods) {
|
||||
this.throwing = Set.of(throwingMethods);
|
||||
}
|
||||
|
||||
private void maybeThrow(String method) {
|
||||
called.add(method);
|
||||
if (throwing.contains(method)) {
|
||||
throw new IllegalStateException("boom: sessions." + method);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
maybeThrow("onTurnComplete");
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean onTurnCompleteWithPostAction(String target) {
|
||||
maybeThrow("onTurnCompleteWithPostAction");
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnFailed(String target) {
|
||||
maybeThrow("onTurnFailed");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnFailed(String target, String reason) {
|
||||
maybeThrow("onTurnFailedWithReason");
|
||||
}
|
||||
}
|
||||
|
||||
private static CompletionResolver newResolver(FakeHerdr herdr, Rendezvous rendezvous) {
|
||||
return new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* A resolver with a controllable clock, and the clock itself, so a test can register a turn
|
||||
* "delivered" at time 0 and then jump the clock past {@link CompletionResolver#MIN_TURN_NANOS}
|
||||
* before resolving it — otherwise {@code onTurnComplete}/{@code onTurnCompleteWithPostAction}
|
||||
* resolve within microseconds of registering in-test, well inside the fleetd#164 floor, and
|
||||
* get classified as a too-fast crash rather than a real completion. Mirrors the {@code
|
||||
* LongSupplier} clock-injection pattern fleetd#164's own tests use.
|
||||
*/
|
||||
private static CompletionResolver newResolverPastTheFloor(FakeHerdr herdr, Rendezvous rendezvous,
|
||||
AtomicLong clock) {
|
||||
LongSupplier nowNanos = clock::get;
|
||||
return new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), nowNanos);
|
||||
}
|
||||
|
||||
@Test
|
||||
void onTurnCompleteResolvesTheWaiterEvenWhenTheSessionHalfThrows() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ real answer\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
AtomicLong clock = new AtomicLong(0);
|
||||
CompletionResolver completion = newResolverPastTheFloor(herdr, rendezvous, clock);
|
||||
var waiter = rendezvous.open("term_a");
|
||||
completion.register("term_a", new TurnToken("term_a", waiter)); // delivered at clock=0
|
||||
clock.set(CompletionResolver.MIN_TURN_NANOS + 1_000_000); // past the too-fast floor
|
||||
|
||||
RecordingSessions sessions = new RecordingSessions("onTurnComplete");
|
||||
TurnListener composed = Fleetd.turnListener(completion, sessions);
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> composed.onTurnComplete("term_a"),
|
||||
"the session half's throw must still escape the composed listener");
|
||||
assertEquals("boom: sessions.onTurnComplete", thrown.getMessage());
|
||||
assertTrue(sessions.called.contains("onTurnComplete"), "the session half must have run");
|
||||
|
||||
// onTurnComplete resolves off-thread (a virtual thread) — wait on the waiter itself,
|
||||
// exactly like MessageServiceTest.completionFallbackIsNeverQueued does.
|
||||
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, resolution.kind(),
|
||||
"the completion resolver must still have resolved the send, despite the session "
|
||||
+ "half throwing");
|
||||
assertEquals("real answer", resolution.text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void onTurnFailedResolvesTheWaiterAsFailedEvenWhenTheSessionHalfThrows() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ crash context\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver completion = newResolver(herdr, rendezvous);
|
||||
var waiter = rendezvous.open("term_b");
|
||||
completion.register("term_b", new TurnToken("term_b", waiter));
|
||||
|
||||
RecordingSessions sessions = new RecordingSessions("onTurnFailed");
|
||||
TurnListener composed = Fleetd.turnListener(completion, sessions);
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> composed.onTurnFailed("term_b"),
|
||||
"the session half's throw must still escape the composed listener");
|
||||
assertEquals("boom: sessions.onTurnFailed", thrown.getMessage());
|
||||
assertTrue(sessions.called.contains("onTurnFailed"), "the session half must have run");
|
||||
|
||||
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(Rendezvous.Kind.FAILED, resolution.kind(),
|
||||
"the completion resolver must still have failed the send, despite the session "
|
||||
+ "half throwing");
|
||||
}
|
||||
|
||||
@Test
|
||||
void onTurnFailedWithReasonResolvesTheWaiterEvenWhenTheSessionHalfThrows() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver completion = newResolver(herdr, rendezvous);
|
||||
var waiter = rendezvous.open("term_c");
|
||||
completion.register("term_c", new TurnToken("term_c", waiter));
|
||||
|
||||
// fleetd #561: the composed listener calls sessions.onTurnFailed(target) — the ONE-arg
|
||||
// overload — for this reason-carrying event too (matching the pre-existing production
|
||||
// behaviour: SessionManager never overrides the two-arg overload either), so the
|
||||
// throwing key here is "onTurnFailed", not a distinct "...WithReason" one.
|
||||
RecordingSessions sessions = new RecordingSessions("onTurnFailed");
|
||||
TurnListener composed = Fleetd.turnListener(completion, sessions);
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> composed.onTurnFailed("term_c", "worker unreachable"),
|
||||
"the session half's throw must still escape the composed listener");
|
||||
assertEquals("boom: sessions.onTurnFailed", thrown.getMessage());
|
||||
assertTrue(sessions.called.contains("onTurnFailed"), "the session half must have run");
|
||||
|
||||
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(Rendezvous.Kind.FAILED, resolution.kind(),
|
||||
"the completion resolver must still have failed the send, despite the session "
|
||||
+ "half throwing");
|
||||
assertEquals("worker unreachable", resolution.text(),
|
||||
"the explicit reason must still reach the resolved send");
|
||||
}
|
||||
|
||||
@Test
|
||||
void onTurnCompleteWithPostActionResolvesTheWaiterEvenWhenTheSessionHalfThrows() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ answer before reset\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
AtomicLong clock = new AtomicLong(0);
|
||||
CompletionResolver completion = newResolverPastTheFloor(herdr, rendezvous, clock);
|
||||
var waiter = rendezvous.open("term_d");
|
||||
completion.register("term_d", new TurnToken("term_d", waiter)); // delivered at clock=0
|
||||
clock.set(CompletionResolver.MIN_TURN_NANOS + 1_000_000); // past the too-fast floor
|
||||
|
||||
RecordingSessions sessions = new RecordingSessions("onTurnCompleteWithPostAction");
|
||||
TurnListener composed = Fleetd.turnListener(completion, sessions);
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> composed.onTurnCompleteWithPostAction("term_d"),
|
||||
"the session half's throw must still escape the composed listener");
|
||||
assertEquals("boom: sessions.onTurnCompleteWithPostAction", thrown.getMessage());
|
||||
assertTrue(sessions.called.contains("onTurnCompleteWithPostAction"), "the session half must have run");
|
||||
|
||||
// resolveBeforePostAction is synchronous by design (it must run before the context reset
|
||||
// can erase the pane) — the waiter is already resolved by the time the throw propagates.
|
||||
assertTrue(waiter.isDone(), "resolveBeforePostAction is synchronous — the send must "
|
||||
+ "already be resolved once the composed call returns (by throwing)");
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
|
||||
assertEquals("answer before reset", waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
/**
|
||||
* The mirror case (comment 17041's acceptance item 2): the completion half throws, and the
|
||||
* session half still ran. The production {@link CompletionResolver} is deliberately defensive
|
||||
* (its scrape reads are wrapped in {@code catch (RuntimeException)}, by the same fail-open
|
||||
* design {@link CompletionResolver#captureBaseline} documents) so it essentially never throws
|
||||
* synchronously in normal operation — forcing it to do so needs a genuinely-real trigger, not
|
||||
* a fabricated one. {@code ConcurrentHashMap.get(null)} is that trigger: passing a {@code null}
|
||||
* target makes {@code onTurnComplete}'s {@code inFlight.get(target)} throw a
|
||||
* {@link NullPointerException} before it ever starts its resolving thread — a real code path,
|
||||
* not a contrived one.
|
||||
*/
|
||||
@Test
|
||||
void sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronously() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver completion = newResolver(herdr, rendezvous);
|
||||
|
||||
RecordingSessions sessions = new RecordingSessions(); // throws from nothing
|
||||
TurnListener composed = Fleetd.turnListener(completion, sessions);
|
||||
|
||||
assertThrows(NullPointerException.class, () -> composed.onTurnComplete(null),
|
||||
"ConcurrentHashMap.get(null) inside CompletionResolver.onTurnComplete must still "
|
||||
+ "escape the composed listener");
|
||||
assertTrue(sessions.called.contains("onTurnComplete"),
|
||||
"the session half must still have run even though the completion half threw first");
|
||||
}
|
||||
|
||||
/**
|
||||
* The mirror case above only exercises {@code onTurnComplete}/{@code bothMustRun}. {@code
|
||||
* onTurnCompleteWithPostAction} is composed through the OTHER helper,
|
||||
* {@code bothMustRunKeepingSecondResult}, and nothing previously asserted that its session half
|
||||
* still runs when its completion half throws — that helper could be reverted to the pre-#561
|
||||
* broken shape (run the session half only if the completion half did not throw) and the suite
|
||||
* would still stay green. Uses the same real, non-fabricated trigger as the test above: a
|
||||
* {@code null} target makes {@code resolveBeforePostAction}'s {@code inFlight.get(target)}
|
||||
* throw a {@link NullPointerException} before {@code resolve} is ever entered.
|
||||
*/
|
||||
@Test
|
||||
void sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver completion = newResolver(herdr, rendezvous);
|
||||
|
||||
RecordingSessions sessions = new RecordingSessions(); // throws from nothing
|
||||
TurnListener composed = Fleetd.turnListener(completion, sessions);
|
||||
|
||||
assertThrows(NullPointerException.class, () -> composed.onTurnCompleteWithPostAction(null),
|
||||
"ConcurrentHashMap.get(null) inside CompletionResolver.resolveBeforePostAction must "
|
||||
+ "still escape the composed listener");
|
||||
assertTrue(sessions.called.contains("onTurnCompleteWithPostAction"),
|
||||
"the session half must still have run even though the completion half threw first");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bothMustRun} must not just run both halves — it must not DROP a second failure when
|
||||
* both halves throw. Forces the completion half to throw (the same real {@code
|
||||
* inFlight.get(null)} NullPointerException trigger used above) while the session half throws a
|
||||
* distinct {@link IllegalStateException}, and asserts the completion half's throwable is what
|
||||
* escapes while the session half's throwable survives as a suppressed exception rather than
|
||||
* being silently discarded.
|
||||
*/
|
||||
@Test
|
||||
void bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver completion = newResolver(herdr, rendezvous);
|
||||
|
||||
RecordingSessions sessions = new RecordingSessions("onTurnComplete");
|
||||
TurnListener composed = Fleetd.turnListener(completion, sessions);
|
||||
|
||||
NullPointerException thrown = assertThrows(NullPointerException.class,
|
||||
() -> composed.onTurnComplete(null),
|
||||
"the completion half's throw (NPE from inFlight.get(null)) must be what escapes");
|
||||
assertTrue(sessions.called.contains("onTurnComplete"), "the session half must still have run");
|
||||
assertEquals(1, thrown.getSuppressed().length,
|
||||
"the session half's distinct failure must be recorded as suppressed, not dropped");
|
||||
assertEquals(IllegalStateException.class, thrown.getSuppressed()[0].getClass());
|
||||
assertEquals("boom: sessions.onTurnComplete", thrown.getSuppressed()[0].getMessage());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.inject.TurnRegistrar;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :505}. {@code Fleetd.main} wires the {@link
|
||||
* dev.ltms.fleet.inject.Injector}'s {@link TurnRegistrar} with {@code completion::register} —
|
||||
* before this ticket that was an inline argument to {@code new Injector(...)}, so nothing could
|
||||
* pin it directly. Measured: replacing it with {@link TurnRegistrar#NOOP} at the call site
|
||||
* compiles with 0 errors and leaves the full suite green, because {@code onDelivered}'s own {@code
|
||||
* captureBaseline} performs the identical {@code inFlight} check-and-put a moment later on the
|
||||
* ordinary delivery path — the two are indistinguishable unless something reads the resolver
|
||||
* between {@code register} and {@code onDelivered}, or {@code onDelivered} never runs at all (the
|
||||
* gap fleetd #556 introduced {@link TurnRegistrar} to close).
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#turnRegistrar} directly — never {@code Injector} or {@code
|
||||
* main} — and drives {@link CompletionResolver} entirely through its public API: {@link
|
||||
* TurnRegistrar#register} followed by {@link CompletionResolver#resolveBeforePostAction}, which
|
||||
* looks up the same {@code inFlight} entry {@code onTurnComplete} would. With the real registrar,
|
||||
* that entry exists and the waiter opened by {@link Rendezvous#open} resolves; with {@link
|
||||
* TurnRegistrar#NOOP} nothing was ever registered, {@code resolveBeforePostAction} finds no
|
||||
* in-flight turn, and the waiter is left exactly as it started — never done.
|
||||
*/
|
||||
class FleetdTurnRegistrarWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.turnRegistrar delegates to the real CompletionResolver, not a no-op")
|
||||
void turnRegistrarDelegatesToCompletionRegister() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ BUILD GREEN: 391 files\n❯ ");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
// fleetd#164: an ever-advancing fake clock stands in for the real time a turn would take
|
||||
// between delivery and resolution, so the MIN_TURN_NANOS "too fast" floor never trips here —
|
||||
// see MessageServiceTest's resolverClock for the same technique.
|
||||
AtomicLong clock = new AtomicLong();
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(),
|
||||
() -> clock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
|
||||
|
||||
TurnRegistrar registrar = Fleetd.turnRegistrar(completion);
|
||||
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
|
||||
registrar.register(T, new TurnToken(T, waiter));
|
||||
|
||||
// Mirrors what Injector.onStatus's confirmed working->idle boundary would trigger via
|
||||
// CompletionResolver.onTurnComplete — resolveBeforePostAction is the public, synchronous
|
||||
// twin of that path and reads the exact same inFlight entry register() must have written.
|
||||
completion.resolveBeforePostAction(T);
|
||||
|
||||
assertTrue(waiter.isDone(),
|
||||
"Fleetd.turnRegistrar(completion) must return completion::register — replacing it "
|
||||
+ "with TurnRegistrar.NOOP at the Fleetd.turnRegistrar call site means this "
|
||||
+ "turn is never registered with CompletionResolver, so resolveBeforePostAction "
|
||||
+ "finds no in-flight turn and this waiter is never resolved");
|
||||
}
|
||||
}
|
||||
@@ -625,6 +625,149 @@ class CompletionResolverTest {
|
||||
assertEquals(Rendezvous.Kind.REPLY, waiterB.getNow(null).kind());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSupersededDoneTurnMustNotEvictItsSuccessorsRegistration() {
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiterA = rendezvous.open("term_a");
|
||||
var turnA = new CompletionResolver.InFlight(waiterA, null);
|
||||
assertTrue(rendezvous.resolve("term_a", "A replied"));
|
||||
rendezvous.close("term_a", waiterA);
|
||||
var waiterB = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
|
||||
|
||||
resolver.resolve("term_a", turnA);
|
||||
|
||||
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
|
||||
"a done turn must not evict B from resolve()'s early return");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSupersededExhaustedTurnMustNotEvictItsSuccessorsRegistration() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ usage limit has been reached\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
target -> Pattern.compile("usage limit has been reached"), ExhaustionSink.none());
|
||||
|
||||
var waiterA = rendezvous.open("term_a");
|
||||
var turnA = new CompletionResolver.InFlight(waiterA, null);
|
||||
rendezvous.close("term_a", waiterA);
|
||||
var waiterB = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
|
||||
|
||||
resolver.resolve("term_a", turnA);
|
||||
|
||||
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
|
||||
"an exhausted turn must not evict B from resolve()'s exhausted branch");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSupersededBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ API Error: 400 invalid request body\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiterA = rendezvous.open("term_a");
|
||||
var turnA = new CompletionResolver.InFlight(waiterA, null);
|
||||
rendezvous.close("term_a", waiterA);
|
||||
var waiterB = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
|
||||
|
||||
resolver.resolve("term_a", turnA);
|
||||
|
||||
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
|
||||
"a backend-error turn must not evict B from resolve()'s error branch");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSupersededRawExhaustedTurnMustNotEvictItsSuccessorsRegistration() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("╭────\nusage limit has been reached");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
target -> Pattern.compile("usage limit has been reached"), ExhaustionSink.none());
|
||||
|
||||
var waiterA = rendezvous.open("term_a");
|
||||
var turnA = new CompletionResolver.InFlight(waiterA, null);
|
||||
rendezvous.close("term_a", waiterA);
|
||||
var waiterB = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
|
||||
|
||||
resolver.resolve("term_a", turnA);
|
||||
|
||||
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
|
||||
"a raw exhausted turn must not evict B from the raw-scrape exhausted branch");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSupersededRawBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("╭────\nAPI Error: 400 invalid request body");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiterA = rendezvous.open("term_a");
|
||||
var turnA = new CompletionResolver.InFlight(waiterA, null);
|
||||
rendezvous.close("term_a", waiterA);
|
||||
var waiterB = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
|
||||
|
||||
resolver.resolve("term_a", turnA);
|
||||
|
||||
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
|
||||
"a raw backend-error turn must not evict B from the raw-scrape error branch");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSupersededDoneFailedTurnMustNotEvictItsSuccessorsRegistration() {
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiterA = rendezvous.open("term_a");
|
||||
var turnA = new CompletionResolver.InFlight(waiterA, null);
|
||||
assertTrue(rendezvous.resolve("term_a", "A replied"));
|
||||
rendezvous.close("term_a", waiterA);
|
||||
var waiterB = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
|
||||
|
||||
resolver.fail("term_a", turnA);
|
||||
|
||||
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
|
||||
"a done failed turn must not evict B from fail()'s early return");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSupersededTooFastBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ API Error: 400 invalid request body\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
long[] clock = {10_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiterA = rendezvous.open("term_a");
|
||||
var turnA = new CompletionResolver.InFlight(waiterA, null, clock[0]);
|
||||
rendezvous.close("term_a", waiterA);
|
||||
var waiterB = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1;
|
||||
|
||||
resolver.resolve("term_a", turnA);
|
||||
|
||||
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
|
||||
"a too-fast backend-error turn must not evict B from failTooFast()");
|
||||
}
|
||||
|
||||
private static void assertSuccessorRegistrationSurvives(CompletionResolver resolver, Rendezvous rendezvous,
|
||||
Object waiterB, String message) {
|
||||
CompletionResolver.InFlight afterA = resolver.inFlight("term_a");
|
||||
assertNotNull(afterA, message + " — a one-arg remove(target) would remove B");
|
||||
assertEquals(waiterB, afterA.waiter(), message + " — the surviving record must belong to B");
|
||||
assertTrue(rendezvous.resolve("term_a", "B replied"), message + " — B must still resolve normally");
|
||||
}
|
||||
|
||||
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -720,10 +720,13 @@ class InjectorTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void anErrorFromSendRemovesTheMessageAndMarksItNotDelivered() {
|
||||
void anErrorFromSendRemovesTheMessageAndMarksItAttempted() {
|
||||
// fleetd #546, acceptance test 1: an Error (not a RuntimeException) escaping the send seam
|
||||
// at Injector.java:383 must still be caught, the message dropped from the queue, and its
|
||||
// state set to NOT_DELIVERED — never left QUEUED at the head.
|
||||
// must still be caught and the message dropped from the queue, never left QUEUED at the
|
||||
// head. fleetd #551 updates what it is marked: an AssertionError from the send call itself
|
||||
// proves nothing about whether the text reached the pane, so it is recorded as ATTEMPTED
|
||||
// (uncertain), not a confident NOT_DELIVERED — see anErrorFromSendMarksItNotDeliveredOnlyForANotFoundCode
|
||||
// below for the one case that still gets NOT_DELIVERED.
|
||||
ErrorOnPrompt throwing = new ErrorOnPrompt(new FakeHerdr());
|
||||
Injector inj = new Injector(new AgentControl(throwing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
|
||||
@@ -734,9 +737,9 @@ class InjectorTest {
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"the delivery's future must surface the send failure");
|
||||
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
|
||||
"the message must be dropped and marked NOT_DELIVERED, not left QUEUED at the "
|
||||
+ "head of the queue");
|
||||
assertEquals(Injector.Cancellation.ATTEMPTED, inj.cancel(delivery),
|
||||
"fleetd #551: the message must be dropped and marked ATTEMPTED (not a confident "
|
||||
+ "NOT_DELIVERED), and never left QUEUED at the head of the queue");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -758,10 +761,14 @@ class InjectorTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void aHerdrExceptionFromSendStillProducesNotDeliveredUnchanged() {
|
||||
void aHerdrExceptionFromSendStillSurfacesButNowReportsAttempted() {
|
||||
// fleetd #546, acceptance test 3 (control): HerdrException extends RuntimeException, so it
|
||||
// was already caught before this ticket's widening. This pins that the ordinary path is
|
||||
// unchanged — still dropped, still NOT_DELIVERED, still surfaced to the caller.
|
||||
// was already caught before #546's widening. The send-failure/surface behaviour is
|
||||
// unchanged by #551 — still dropped, still surfaced to the caller — but fleetd #551
|
||||
// deliberately changes WHAT it is marked: "send_failed" is not a herdr `*_not_found` code,
|
||||
// so this is exactly the response-half case #551 exists to fix. See
|
||||
// aHerdrExceptionAfterThePasteIsNeverRecordedAsConfidentlyNotDelivered for the acceptance
|
||||
// test this ticket was filed for.
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
@@ -771,9 +778,90 @@ class InjectorTest {
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"a HerdrException at the send seam must still surface to the caller, unchanged by "
|
||||
+ "fleetd #546's widening");
|
||||
assertEquals(Injector.Cancellation.ATTEMPTED, inj.cancel(delivery),
|
||||
"fleetd #551: a non-*_not_found HerdrException must now report ATTEMPTED, not a "
|
||||
+ "confident NOT_DELIVERED");
|
||||
}
|
||||
|
||||
// --- fleetd #551: the injector records the delivery attempt BEFORE the irreversible send, not
|
||||
// after, so a failure in the send's response window can never be written down as a confident
|
||||
// NOT_DELIVERED for text that may already be sitting in the worker's pane. ---
|
||||
|
||||
@Test
|
||||
void aHerdrExceptionAfterThePasteIsNeverRecordedAsConfidentlyNotDelivered() {
|
||||
// fleetd #551, the ticket's acceptance test 1, and comment 17037's added test ("a test
|
||||
// pinning that a caller cannot receive a confident NOT_DELIVERED for a message that reached
|
||||
// the paste"): a HerdrException thrown from the response half of the send seam — herdr
|
||||
// replied, with an error, so it definitely processed the request — must not leave a record
|
||||
// claiming the text was never delivered. Red before the fix: on main at ba2f4d1 this
|
||||
// asserted (and got) Injector.Cancellation.NOT_DELIVERED.
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"the caller must still see the send failure");
|
||||
assertEquals(Injector.Cancellation.ATTEMPTED, inj.cancel(delivery),
|
||||
"fleetd #551: a HerdrException from the response half of the send seam must not be "
|
||||
+ "recorded as a confident NOT_DELIVERED — the text may already be sitting "
|
||||
+ "in the pane");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ordinarySuccessStillReportsDeliveredExactlyOnce() {
|
||||
// fleetd #551, the ticket's acceptance test 2: the record-before-send reorder must not
|
||||
// change the ordinary success path — exactly one send, and the caller-facing Cancellation
|
||||
// for it is still DELIVERED, never left at the new ATTEMPTED value.
|
||||
Injector.Delivery delivery = injector.enqueue(T, "hello", TestTurnTokens.inert(T));
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertEquals(List.of("hello"), sent(), "the text must be sent exactly once");
|
||||
assertTrue(delivery.completion().isDone() && !delivery.completion().isCompletedExceptionally(),
|
||||
"the ordinary success path must still complete normally");
|
||||
assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(delivery),
|
||||
"fleetd #551: the ordinary success path must still report DELIVERED, unaffected by "
|
||||
+ "the record-before-send reorder");
|
||||
}
|
||||
|
||||
@Test
|
||||
void transportDownStillSurfacesTheErrorAndDropsTheEntry() {
|
||||
// fleetd #551, the ticket's acceptance test 3: the ordinary transport-down path (herdr
|
||||
// unreachable — a plain HerdrException with no error code) must be unaffected by the #551
|
||||
// reorder — the failure still surfaces to the caller, and the entry is not left sitting in
|
||||
// the queue for a second onStatus round to resend.
|
||||
FakeHerdr unreachable = new FakeHerdr().healthy(false);
|
||||
Injector inj = new Injector(new AgentControl(unreachable));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"a transport-down failure must still surface to the caller");
|
||||
assertTrue(inj.activeTargets().isEmpty(),
|
||||
"the entry must not be left queued for a second round to resend");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNotFoundCodeFromSendStillReportsAConfidentNotDelivered() {
|
||||
// fleetd #551 (comment 17037, point 3): a herdr `*_not_found` error means the target
|
||||
// pane/agent does not exist at all, so nothing could have been pasted anywhere — the
|
||||
// codebase already treats this family as a confirmed absence everywhere else (StatusPoller,
|
||||
// AgentControl's own retry). #551's new ATTEMPTED state must not swallow this case: it stays
|
||||
// a confident NOT_DELIVERED.
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("agent_not_found");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"an agent_not_found failure must still surface to the caller");
|
||||
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
|
||||
"a HerdrException must still be dropped and marked NOT_DELIVERED, unchanged by the "
|
||||
+ "wider Throwable catch");
|
||||
"fleetd #551: a herdr *_not_found error is a confirmed absence, not merely "
|
||||
+ "inconclusive — it must stay NOT_DELIVERED, not the new ATTEMPTED");
|
||||
}
|
||||
|
||||
// --- fleetd #553: a throwable from any listener callback in onStatus must not skip the
|
||||
|
||||
@@ -7,7 +7,9 @@ import dev.ltms.fleet.auth.Role;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.LoopWatchdog;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
@@ -57,9 +59,10 @@ class FleetMcpTest {
|
||||
|
||||
private final FakeHerdr herdr = new FakeHerdr();
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private final Injector injector = new Injector(agents);
|
||||
private final Rendezvous rendezvous = new Rendezvous();
|
||||
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
private final MessageService messages = new MessageService(agents, new Injector(agents), rendezvous, inbox);
|
||||
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
@@ -298,6 +301,35 @@ class FleetMcpTest {
|
||||
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #571 (ticket CORRECTION 5): {@code formatReply}'s {@code TIMED_OUT_UNCONFIRMED} arm is
|
||||
* the one message whose whole job is to stop a caller retrying a delivery that may already have
|
||||
* arrived. Pin that its wording is actually distinct from the queued/working arm's retry
|
||||
* invitation — a mutation that swapped this arm's text for that one still passed every other
|
||||
* test in this suite, because nothing asserted the specific wording.
|
||||
*/
|
||||
@Test
|
||||
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
|
||||
herdr.agentSendFailsWith("send_failed");
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of()));
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt -> ATTEMPTED
|
||||
|
||||
McpSchema.CallToolResult res = send.get(5, TimeUnit.SECONDS);
|
||||
String text = textOf(res);
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
|
||||
assertTrue(text.contains("delivery unconfirmed"), "got: " + text);
|
||||
assertFalse(text.contains("retry or poll status"),
|
||||
"an unconfirmed delivery must not carry the queued/working arm's retry invitation — "
|
||||
+ "a resend here can double-deliver the same brief: got " + text);
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendRejectsMissingArgs() {
|
||||
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError());
|
||||
@@ -1053,6 +1085,36 @@ class FleetMcpTest {
|
||||
assertFalse(out.contains("quarantinedForSeconds"), out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void loopHealthReportsStalledStatusPoller() {
|
||||
String out = loopHealth(LoopWatchdog.State.STALLED, LoopWatchdog.State.RUNNING);
|
||||
assertTrue(out.contains("\"statusPoller\":\"STALLED\""),
|
||||
"fleet_list must report a stalled StatusPoller: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void loopHealthReportsStoppedSessionReaperAsStopped() {
|
||||
String out = loopHealth(LoopWatchdog.State.RUNNING, LoopWatchdog.State.STOPPED);
|
||||
assertTrue(out.contains("\"sessionReaper\":\"STOPPED\""),
|
||||
"fleet_list must report a deliberately stopped SessionReaper as STOPPED, not an alarm: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void loopHealthReportsRunningStatusPoller() {
|
||||
String out = loopHealth(LoopWatchdog.State.RUNNING, LoopWatchdog.State.STOPPED);
|
||||
assertTrue(out.contains("\"statusPoller\":\"RUNNING\""),
|
||||
"fleet_list must report a running StatusPoller: " + out);
|
||||
}
|
||||
|
||||
private static String loopHealth(LoopWatchdog.State statusPoller, LoopWatchdog.State sessionReaper) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
return textOf(FleetMcp.listFleet(workerService(herdr, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
new SessionManager(workerService(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"))), null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
new FleetMcp.LoopHealthSource(() -> statusPoller, () -> sessionReaper),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), ""));
|
||||
}
|
||||
|
||||
@Test
|
||||
void capacityIncludesConfiguredProfileWithoutMembers() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
|
||||
@@ -8,7 +8,10 @@ import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.Timeout;
|
||||
|
||||
import java.lang.reflect.InvocationHandler;
|
||||
import java.lang.reflect.Method;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.io.IOException;
|
||||
import java.util.Map;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
@@ -20,6 +23,7 @@ import java.util.concurrent.atomic.AtomicReference;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* CB-528 follow-up: {@link AmqpReplyInbox#failPendingPublishesOnRecovery()} must not fail a publish
|
||||
@@ -197,11 +201,78 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
+ elapsedMillis.get() + "ms");
|
||||
}
|
||||
|
||||
@Test
|
||||
void publishIOExceptionRemovesThePendingMessageId() throws Exception {
|
||||
AtomicLong seqCounter = new AtomicLong();
|
||||
Channel failing = fakeChannel(seqCounter, new CopyOnWriteArrayList<>(), new CopyOnWriteArrayList<>(),
|
||||
new AtomicReference<>(), new AtomicReference<>(), true);
|
||||
AmqpReplyInbox inbox = new AmqpReplyInbox(fakeConnection(failing, failing), AmqpReplyInbox.DEFAULT_PREFETCH);
|
||||
|
||||
try {
|
||||
org.junit.jupiter.api.Assertions.assertThrows(IllegalStateException.class,
|
||||
() -> inbox.publish("worker", "catch", "body"));
|
||||
assertEquals(0, pendingByMsgId(inbox).size(), "publish IOException must remove its msgId entry");
|
||||
} finally {
|
||||
inbox.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void interruptedPublishRemovesThePendingMessageIdInFinally() throws Exception {
|
||||
InboxFixture fixture = new InboxFixture();
|
||||
Thread publish = fixture.startPublish("finally");
|
||||
fixture.awaitPublished("finally");
|
||||
publish.interrupt();
|
||||
publish.join(5_000);
|
||||
|
||||
assertEquals(0, pendingByMsgId(fixture.inbox).size(), "publish finally must remove its msgId entry");
|
||||
fixture.inbox.close();
|
||||
}
|
||||
|
||||
@Test
|
||||
void confirmResolutionRemovesThePendingMessageId() throws Exception {
|
||||
AmqpReplyInbox inbox = new InboxFixture().inbox;
|
||||
try {
|
||||
seedPending(inbox, 1, "confirm");
|
||||
invoke(inbox, "resolveConfirm", new Class<?>[] {long.class, boolean.class, boolean.class}, 1L, false, true);
|
||||
assertEquals(0, pendingByMsgId(inbox).size(), "confirm resolution must remove its msgId entry");
|
||||
} finally {
|
||||
inbox.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void recoverySweepRemovesThePendingMessageId() throws Exception {
|
||||
AmqpReplyInbox inbox = new InboxFixture().inbox;
|
||||
try {
|
||||
seedPending(inbox, 1, "recovery");
|
||||
inbox.failPendingPublishesOnRecovery();
|
||||
assertEquals(0, pendingByMsgId(inbox).size(), "recovery sweep must remove its msgId entry");
|
||||
} finally {
|
||||
inbox.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void closeRemovesThePendingMessageId() throws Exception {
|
||||
AmqpReplyInbox inbox = new InboxFixture().inbox;
|
||||
seedPending(inbox, 1, "close");
|
||||
inbox.close();
|
||||
|
||||
assertEquals(0, pendingByMsgId(inbox).size(), "close must remove its msgId entry");
|
||||
}
|
||||
|
||||
/** A {@link Proxy}-backed {@link Channel}: only the calls {@link AmqpReplyInbox} actually makes
|
||||
* are meaningfully implemented; everything else returns a harmless default. */
|
||||
private static Channel fakeChannel(AtomicLong seqCounter, List<Long> seqOrder, List<String> msgIdOrder,
|
||||
AtomicReference<ConfirmCallback> ackCallback,
|
||||
AtomicReference<ConfirmCallback> nackCallback) {
|
||||
AtomicReference<ConfirmCallback> ackCallback,
|
||||
AtomicReference<ConfirmCallback> nackCallback) {
|
||||
return fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback, false);
|
||||
}
|
||||
|
||||
private static Channel fakeChannel(AtomicLong seqCounter, List<Long> seqOrder, List<String> msgIdOrder,
|
||||
AtomicReference<ConfirmCallback> ackCallback,
|
||||
AtomicReference<ConfirmCallback> nackCallback, boolean failPublish) {
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
String name = method.getName();
|
||||
if (name.equals("getNextPublishSeqNo")) {
|
||||
@@ -210,6 +281,9 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
return value;
|
||||
}
|
||||
if (name.equals("basicPublish")) {
|
||||
if (failPublish) {
|
||||
throw new IOException("test publish failure");
|
||||
}
|
||||
AMQP.BasicProperties props = (AMQP.BasicProperties) args[3];
|
||||
msgIdOrder.add(props.getMessageId());
|
||||
return null;
|
||||
@@ -285,4 +359,60 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, Object> pendingByMsgId(AmqpReplyInbox inbox) throws Exception {
|
||||
var field = AmqpReplyInbox.class.getDeclaredField("pendingByMsgId");
|
||||
field.setAccessible(true);
|
||||
return (Map<String, Object>) field.get(inbox);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static void seedPending(AmqpReplyInbox inbox, long seq, String msgId) throws Exception {
|
||||
Class<?> pendingType = Class.forName(AmqpReplyInbox.class.getName() + "$Pending");
|
||||
var constructor = pendingType.getDeclaredConstructor(String.class);
|
||||
constructor.setAccessible(true);
|
||||
Object pending = constructor.newInstance(msgId);
|
||||
var seqField = AmqpReplyInbox.class.getDeclaredField("pendingBySeq");
|
||||
seqField.setAccessible(true);
|
||||
((Map<Long, Object>) seqField.get(inbox)).put(seq, pending);
|
||||
pendingByMsgId(inbox).put(msgId, pending);
|
||||
}
|
||||
|
||||
private static void invoke(AmqpReplyInbox inbox, String name, Class<?>[] types, Object... args) throws Exception {
|
||||
Method method = AmqpReplyInbox.class.getDeclaredMethod(name, types);
|
||||
method.setAccessible(true);
|
||||
method.invoke(inbox, args);
|
||||
}
|
||||
|
||||
private static final class InboxFixture {
|
||||
final AtomicLong seqCounter = new AtomicLong();
|
||||
final List<Long> seqOrder = new CopyOnWriteArrayList<>();
|
||||
final List<String> msgIdOrder = new CopyOnWriteArrayList<>();
|
||||
final AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
|
||||
final AtomicReference<ConfirmCallback> nackCallback = new AtomicReference<>();
|
||||
final AmqpReplyInbox inbox = new AmqpReplyInbox(
|
||||
fakeConnection(fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback),
|
||||
fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback)),
|
||||
AmqpReplyInbox.DEFAULT_PREFETCH);
|
||||
|
||||
Thread startPublish(String msgId) {
|
||||
Thread thread = Thread.ofVirtual().start(() -> {
|
||||
try {
|
||||
inbox.publish("worker", msgId, "body");
|
||||
} catch (IllegalStateException ignored) {
|
||||
// Interrupting the confirm wait is the path under test.
|
||||
}
|
||||
});
|
||||
return thread;
|
||||
}
|
||||
|
||||
void awaitPublished(String msgId) throws InterruptedException {
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
|
||||
while (!msgIdOrder.contains(msgId) && System.nanoTime() < deadline) {
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertTrue(msgIdOrder.contains(msgId), "publish did not register " + msgId);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,13 +12,18 @@ import org.testcontainers.junit.jupiter.Testcontainers;
|
||||
import org.testcontainers.utility.DockerImageName;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.lang.reflect.InvocationHandler;
|
||||
import java.lang.reflect.Method;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
@@ -207,6 +212,31 @@ class LeadMailboxTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link LeadMailbox#inspect} opens a third channel after the mailbox's consume and publish
|
||||
* channels. Limit this connection to three channels, then require a replacement channel after
|
||||
* the successful inspect. If inspect leaves its probe open, the broker refuses that replacement.
|
||||
*/
|
||||
@Test
|
||||
void inspectClosesItsSuccessfulProbeChannel() throws Exception {
|
||||
var factory = LeadMailbox.connectionFactory(uri());
|
||||
factory.setRequestedChannelMax(3);
|
||||
Connection connection = factory.newConnection();
|
||||
try (LeadMailbox mailbox = new LeadMailbox(connection, coordId("lead-inspect-probe-close"))) {
|
||||
LeadChannel.MailboxState state = mailbox.inspect(mailbox.selfCoordId());
|
||||
assertTrue(state.exists(), "the owned mailbox must be found before checking the probe channel");
|
||||
|
||||
Channel replacement = connection.createChannel();
|
||||
assertNotNull(replacement,
|
||||
"inspect must close its successful probe channel; the replacement channel was null");
|
||||
try {
|
||||
assertTrue(replacement.isOpen(), "the replacement channel must be open after inspect returns");
|
||||
} finally {
|
||||
replacement.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #440: {@code heldDurable()} must be derived from what {@link LeadMailbox#own} actually
|
||||
* did against the real broker — a durable queue declare plus a manual-ack consumer — not a
|
||||
@@ -347,6 +377,72 @@ class LeadMailboxTest {
|
||||
() -> "expected AlreadyClosedException, got: " + thrown);
|
||||
}
|
||||
|
||||
@Test
|
||||
void publishIOExceptionRemovesThePendingMessageId() throws Exception {
|
||||
LeadMailbox mailbox = newMailbox(true);
|
||||
try {
|
||||
assertThrows(IllegalStateException.class,
|
||||
() -> mailbox.publish("target", new LeadMessage("catch", "from", "target", "body")));
|
||||
assertEquals(0, pendingByMsgId(mailbox).size(), "publish IOException must remove its msgId entry");
|
||||
} finally {
|
||||
mailbox.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void interruptedPublishRemovesThePendingMessageIdInFinally() throws Exception {
|
||||
LeadMailbox mailbox = newMailbox(false);
|
||||
Thread publish = Thread.ofVirtual().start(() -> {
|
||||
try {
|
||||
mailbox.publish("target", new LeadMessage("finally", "from", "target", "body"));
|
||||
} catch (IllegalStateException ignored) {
|
||||
// Interrupting the confirm wait is the path under test.
|
||||
}
|
||||
});
|
||||
awaitPending(mailbox, "finally");
|
||||
publish.interrupt();
|
||||
publish.join(5_000);
|
||||
|
||||
try {
|
||||
assertEquals(0, pendingByMsgId(mailbox).size(), "publish finally must remove its msgId entry");
|
||||
} finally {
|
||||
mailbox.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void confirmResolutionRemovesThePendingMessageId() throws Exception {
|
||||
LeadMailbox mailbox = newMailbox(false);
|
||||
try {
|
||||
seedPending(mailbox, 1, "confirm");
|
||||
invoke(mailbox, "resolveConfirm", new Class<?>[] {long.class, boolean.class, boolean.class}, 1L, false, true);
|
||||
assertEquals(0, pendingByMsgId(mailbox).size(), "confirm resolution must remove its msgId entry");
|
||||
} finally {
|
||||
mailbox.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void recoverySweepRemovesThePendingMessageId() throws Exception {
|
||||
LeadMailbox mailbox = newMailbox(false);
|
||||
try {
|
||||
seedPending(mailbox, 1, "recovery");
|
||||
mailbox.failPendingPublishesOnRecovery();
|
||||
assertEquals(0, pendingByMsgId(mailbox).size(), "recovery sweep must remove its msgId entry");
|
||||
} finally {
|
||||
mailbox.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void closeRemovesThePendingMessageId() throws Exception {
|
||||
LeadMailbox mailbox = newMailbox(false);
|
||||
seedPending(mailbox, 1, "close");
|
||||
mailbox.close();
|
||||
|
||||
assertEquals(0, pendingByMsgId(mailbox).size(), "close must remove its msgId entry");
|
||||
}
|
||||
|
||||
/** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */
|
||||
@SuppressWarnings("BusyWait")
|
||||
private static List<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
|
||||
@@ -372,4 +468,108 @@ class LeadMailboxTest {
|
||||
}
|
||||
return state;
|
||||
}
|
||||
|
||||
private static LeadMailbox newMailbox(boolean failPublish) {
|
||||
AtomicLong sequence = new AtomicLong();
|
||||
Channel consume = fakeChannel(sequence, false);
|
||||
Channel publish = fakeChannel(sequence, failPublish);
|
||||
return new LeadMailbox(fakeConnection(consume, publish), "self");
|
||||
}
|
||||
|
||||
private static Channel fakeChannel(AtomicLong sequence, boolean failPublish) {
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
if (method.getName().equals("getNextPublishSeqNo")) {
|
||||
return sequence.incrementAndGet();
|
||||
}
|
||||
if (method.getName().equals("basicPublish") && failPublish) {
|
||||
throw new IOException("test publish failure");
|
||||
}
|
||||
if (method.getName().equals("equals")) {
|
||||
return proxy == args[0];
|
||||
}
|
||||
if (method.getName().equals("hashCode")) {
|
||||
return System.identityHashCode(proxy);
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Channel) Proxy.newProxyInstance(LeadMailboxTest.class.getClassLoader(), new Class<?>[] {Channel.class}, handler);
|
||||
}
|
||||
|
||||
private static Connection fakeConnection(Channel first, Channel second) {
|
||||
AtomicLong calls = new AtomicLong();
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
if (method.getName().equals("createChannel") && (args == null || args.length == 0)) {
|
||||
return calls.getAndIncrement() == 0 ? first : second;
|
||||
}
|
||||
if (method.getName().equals("equals")) {
|
||||
return proxy == args[0];
|
||||
}
|
||||
if (method.getName().equals("hashCode")) {
|
||||
return System.identityHashCode(proxy);
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Connection) Proxy.newProxyInstance(LeadMailboxTest.class.getClassLoader(), new Class<?>[] {Connection.class}, handler);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, Object> pendingByMsgId(LeadMailbox mailbox) throws Exception {
|
||||
var field = LeadMailbox.class.getDeclaredField("pendingByMsgId");
|
||||
field.setAccessible(true);
|
||||
return (Map<String, Object>) field.get(mailbox);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static void seedPending(LeadMailbox mailbox, long seq, String msgId) throws Exception {
|
||||
Class<?> pendingType = Class.forName(LeadMailbox.class.getName() + "$Pending");
|
||||
var constructor = pendingType.getDeclaredConstructor(String.class);
|
||||
constructor.setAccessible(true);
|
||||
Object pending = constructor.newInstance(msgId);
|
||||
var seqField = LeadMailbox.class.getDeclaredField("pendingBySeq");
|
||||
seqField.setAccessible(true);
|
||||
((Map<Long, Object>) seqField.get(mailbox)).put(seq, pending);
|
||||
pendingByMsgId(mailbox).put(msgId, pending);
|
||||
}
|
||||
|
||||
private static void invoke(LeadMailbox mailbox, String name, Class<?>[] types, Object... args) throws Exception {
|
||||
Method method = LeadMailbox.class.getDeclaredMethod(name, types);
|
||||
method.setAccessible(true);
|
||||
method.invoke(mailbox, args);
|
||||
}
|
||||
|
||||
private static void awaitPending(LeadMailbox mailbox, String msgId) throws Exception {
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
|
||||
while (!pendingByMsgId(mailbox).containsKey(msgId) && System.nanoTime() < deadline) {
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertTrue(pendingByMsgId(mailbox).containsKey(msgId), "publish did not register " + msgId);
|
||||
}
|
||||
|
||||
private static Object defaultValue(Class<?> type) {
|
||||
if (!type.isPrimitive() || type == void.class) {
|
||||
return null;
|
||||
}
|
||||
if (type == boolean.class) {
|
||||
return Boolean.FALSE;
|
||||
}
|
||||
if (type == long.class) {
|
||||
return 0L;
|
||||
}
|
||||
if (type == short.class) {
|
||||
return (short) 0;
|
||||
}
|
||||
if (type == byte.class) {
|
||||
return (byte) 0;
|
||||
}
|
||||
if (type == char.class) {
|
||||
return (char) 0;
|
||||
}
|
||||
if (type == double.class) {
|
||||
return 0.0d;
|
||||
}
|
||||
if (type == float.class) {
|
||||
return 0.0f;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -506,6 +506,52 @@ class MessageServiceTest {
|
||||
"an answer to a turn that never existed (or already lapsed) is stale, not a hang");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #575. {@code answer()}'s STALE_TURN return for "the ask lapsed between the lookup and
|
||||
* the unblock" ({@code rendezvous.answerAsk(turnId, content)} returning {@code false} even though
|
||||
* this call's own {@code rendezvous.askSession(turnId)} check at the top saw the ask as open) used
|
||||
* to be covered only by a hand-rolled copy of the cleanup pair its own finally already runs, not
|
||||
* the finally itself — a structural gap, closed by widening the try up to cover the registration
|
||||
* above it. This test pins the exact race with {@link
|
||||
* MessageService#setAnswerAskLapseRaceHookForTest}, which fires right before this call's own
|
||||
* {@code rendezvous.answerAsk} call and completes the same turnId's ask directly — reproducing
|
||||
* what a second, concurrent {@code answer()} winning that race could otherwise only do by timing
|
||||
* luck. Proves the fix changed no behaviour on this path: it still returns {@code STALE_TURN},
|
||||
* and this call's own forward waiter is still closed exactly once (not left open, and not closed
|
||||
* twice — there is now only one cleanup site left to run).
|
||||
*/
|
||||
@Test
|
||||
void answerLosingTheRaceToAnAlreadyAnsweredAskStillReturnsStaleTurnAndCleansUpOnce() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
String turnId = asking.turnId();
|
||||
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
|
||||
|
||||
messages.setAnswerAskLapseRaceHookForTest(() -> rendezvous.answerAsk(turnId, "raced in first"));
|
||||
try {
|
||||
MessageService.Reply r = messages.answer(turnId, "too late", 500);
|
||||
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
|
||||
"an ask already answered by the race must be seen as lapsed, not double-delivered");
|
||||
} finally {
|
||||
messages.setAnswerAskLapseRaceHookForTest(null);
|
||||
}
|
||||
|
||||
// Cleanup ran exactly once: the forward waiter THIS call opened is closed, not leaked.
|
||||
assertNull(rendezvous.currentWaiter(T),
|
||||
"the forward waiter this answer() call opened must be closed after a STALE_TURN return");
|
||||
|
||||
// The worker's own ask() call, unblocked by the hook's direct answerAsk, still completes
|
||||
// normally — the race this test simulates does not strand it.
|
||||
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome());
|
||||
assertEquals("raced in first", a.answer());
|
||||
}
|
||||
|
||||
// --- timeout, answer, poll, and lock-contention edges ----------------------------------
|
||||
|
||||
@Test
|
||||
@@ -555,6 +601,30 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #571 (the acceptance test the ticket was filed for). The worker is idle so the injector
|
||||
* attempts delivery, but the {@code agent.prompt} call itself fails with a herdr error that is
|
||||
* not a confirmed absence (not a {@code *_not_found} code) — {@link Injector} marks the Pending
|
||||
* {@code ATTEMPTED} (fleetd #551), meaning the call was made and whether it reached the pane is
|
||||
* unknown. Before this fix, {@code send}'s {@code TimeoutException} branch collapsed
|
||||
* {@code ATTEMPTED} into {@code TIMED_OUT_QUEUED} — a promise that the message will never arrive,
|
||||
* which may already be false: {@code agent.prompt} pastes and submits in one call.
|
||||
*/
|
||||
@Test
|
||||
void sendTimesOutWithAttemptedDeliveryReportsUnconfirmedNotQueued() throws Exception {
|
||||
herdr.agentSendFailsWith("send_failed");
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt → ATTEMPTED
|
||||
|
||||
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_UNCONFIRMED, r.outcome(),
|
||||
"an ATTEMPTED delivery must not collapse into TIMED_OUT_QUEUED — the message may "
|
||||
+ "already have arrived in full, and TIMED_OUT_QUEUED promises it never will");
|
||||
assertNull(r.text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
|
||||
@@ -8,6 +8,7 @@ import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.LoopWatchdog;
|
||||
import dev.ltms.fleet.inject.StatusPoller;
|
||||
import dev.ltms.fleet.inject.MemberPresence;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
@@ -165,6 +166,29 @@ class FleetAppTest {
|
||||
assertEquals("degraded", mapper.readTree(res.body()).get("status").asText());
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzKeepsOkStatusAndReportsLoopHealthInItsBody() {
|
||||
FleetMcp.LoopHealthSource loops = new FleetMcp.LoopHealthSource(
|
||||
() -> LoopWatchdog.State.RUNNING, () -> LoopWatchdog.State.STOPPED);
|
||||
|
||||
FleetApp.HealthzResponse ok = FleetApp.healthzResponse(new FakeHerdr(), new FakeHerdr(), loops);
|
||||
assertEquals(200, ok.status(), "a healthy herdr must keep /healthz at 200 regardless of loop states");
|
||||
assertEquals(Map.of("statusPoller", "RUNNING", "sessionReaper", "STOPPED"), ok.body().get("loopHealth"),
|
||||
"the /healthz body must report each loop state without making STOPPED an alarm");
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzKeepsDegradedStatusAndReportsLoopHealthInItsBody() {
|
||||
FleetMcp.LoopHealthSource loops = new FleetMcp.LoopHealthSource(
|
||||
() -> LoopWatchdog.State.RUNNING, () -> LoopWatchdog.State.STOPPED);
|
||||
FleetApp.HealthzResponse degraded = FleetApp.healthzResponse(new FakeHerdr().healthy(false),
|
||||
new FakeHerdr(), loops);
|
||||
assertEquals(503, degraded.status(),
|
||||
"an unreachable herdr must keep /healthz at 503 regardless of loop states");
|
||||
assertEquals(Map.of("statusPoller", "RUNNING", "sessionReaper", "STOPPED"),
|
||||
degraded.body().get("loopHealth"), "the degraded /healthz body must retain loop states");
|
||||
}
|
||||
|
||||
@Test
|
||||
void sessionsMapsWorkspaceList() throws Exception {
|
||||
int port = startHealthy();
|
||||
@@ -623,6 +647,28 @@ class FleetAppTest {
|
||||
assertTrue(herdr.called("agent.prompt"), "message was injected");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #571: the worker is idle, so the poller attempts delivery, but the {@code agent.prompt}
|
||||
* call itself fails with a herdr error that is not a confirmed absence — {@link
|
||||
* dev.ltms.fleet.inject.Injector} marks this {@code ATTEMPTED}, meaning the call was made and
|
||||
* whether it reached the pane is unknown. {@code writeReply}'s default arm must map this to its
|
||||
* own {@code "unconfirmed"} status, not silently fall through to {@code "done"} (which would
|
||||
* claim the delegation completed) nor collapse into {@code "queued"} (which would claim the
|
||||
* message will never arrive, when it may already be sitting in the pane).
|
||||
*/
|
||||
@Test
|
||||
void messageTimesOutUnconfirmedWhenDeliveryAttemptFails() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("send_failed");
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":250}");
|
||||
assertEquals(202, res.statusCode());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("unconfirmed", body.get("status").asText(),
|
||||
"an ATTEMPTED delivery must report its own status, not \"queued\" or \"done\"");
|
||||
assertTrue(herdr.called("agent.prompt"), "delivery must have been attempted");
|
||||
}
|
||||
|
||||
@Test
|
||||
void messageRejectsBlankContent() throws Exception {
|
||||
int port = startHealthy();
|
||||
|
||||
Reference in New Issue
Block a user