Compare commits

...

13 Commits

Author SHA1 Message Date
Dai Ha 6cccd458d4 #589 Group 3: wiring-test the 5 sites below line 500 in Fleetd.main()
CI / shell-tests (pull_request) Successful in 11s
CI / contract (pull_request) Successful in 1m26s
CI / build (pull_request) Successful in 2m11s
Extracts the inline lambdas/method references at the 5 assigned wiring
sites into named package-private factories on Fleetd, following the
FleetdLoopHealthSourceWiringTest pattern from #584:

- turnRegistrar(CompletionResolver) — was completion::register (Injector)
- healthFailTarget(MessageService) — was messages::abandon (FleetHealthMonitor)
- releaseCleanup(MessageService, ReplyInbox, PrimaryRegistry) — was the
  inline sessions.onRelease(detail -> {...}) cleanup lambda
- replyInboxOpener() — was AmqpReplyInbox::open passed to selectReplyInbox
- leadMailboxOpener() — was LeadMailbox::open passed to openLeadMailbox

Each factory has a new runtime test (not source-text) that drives real
collaborators through public APIs: MessageService.poll(ticket).phase(),
InMemoryReplyInbox.peek(), PrimaryRegistry.nudgeTargetFor(), and the
opener tests connect to a guaranteed-closed local port to prove a real
network attempt vs. an inert stub.

releaseCleanup was done first per the brief: MessageService.abandon's
javadoc documents that losing this cleanup leaves a torn-down worker's
rendezvous waiter open forever.

Tests: 1789 -> 1794 (+5), 0 failures, 0 errors. mvn -q -o test exit 0,
no BUILD FAILURE, no piped exit status. Each new test verified RED on
the inert form named in the ticket, and GREEN after reformatting the
call across lines and extracting the argument into a local/factory.
2026-09-19 15:24:11 +07:00
Dai Ha a639969a9a CLAUDE.md: a blocked lead consults architects, not the operator
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 1m32s
CI / build (push) Successful in 1m42s
The operator set this rule on 2026-09-19: when a decision blocks a lead, it
consults one or more architect members, who are authorized to agree on one
decision and unblock. The operator is not asked. Escalation stays open only
for things outside the fleet's authority -- money, credentials, or a promise
made to someone else.

The paragraph also carries the reason the ticket record is mandatory rather
than optional. The operator's old notification channel was the block itself:
work stopped, so they found out. Taking the operator out of the loop removes
that signal with it, so the decision goes on the ticket, which reaches them
whether or not they are at a terminal when it is made.

The rule has a second half aimed at architects, which lives in
fleet.charters.architect and is applied per daemon -- filed as #591, because
a charter can never reach a lead and charters do not travel between hosts.
The missing notification event is #592.

Block verified byte-identical with wiki 7-Use-Cases.md at 1d9bd1b.
2026-09-19 14:57:04 +07:00
Dai Ha 17c3a69c57 docs: CB-591 gateway page had two claims that went stale
CI / shell-tests (push) Successful in 5s
CI / contract (push) Successful in 1m23s
CI / build (push) Failing after 1m37s
The gateway's chat model is served under the stable alias `acoder`, and the
model behind that alias changed on 2026-08-28 — it is Qwen3.8-27B now, not
DeepSeek-V4-Flash. The old name is still served, so nothing broke, but it
names a model this is not.

Two claims on the page were wrong as a result, and both were written as
current facts rather than dated measurements:

- `/v1/models` returns exactly `["deepseek-v4-flash"]` — it returns 6 ids
  now. This sat under a heading saying it needs no re-testing.
- the status banner said `local` and `gx` are both at `weight: 100` —
  `local` is at 0.

Measured today against the live gateway: /v1/models returns acoder,
qwen3.8-27b-nvfp4, deepseek-v4-flash and three embedding names; a completion
sent as `deepseek-v4-flash` comes back reporting `"model": "acoder"`, which
is the alias in plain sight. /v1/deployment reports generation
2026-08-28-qwen3.8-27b-nvfp4.

§2 and §3 are left alone. They are the August plan, and rewriting them would
destroy the record of the migration.

fleetd.yaml moved to `acoder` in the same change. It is not tracked here.
2026-09-13 07:01:15 +07:00
ltms 49a5875586 Merge #583: fleetd #582 — assert pending message-id cleanup at every publish cleanup site
CI / shell-tests (push) Successful in 9s
CI / contract (push) Successful in 48s
CI / build (push) Failing after 2m1s
All ten assertions proven live: six by the implementer, the last four by the lead.
One contract build with four deleted removal lines produced exactly four named failures,
one per site, with the total unchanged at 1825.
2026-09-12 15:52:25 +02:00
Dai Ha 634d33b50b Merge worker/562-loop-health-wiring-test-99611c-5
CI / shell-tests (push) Successful in 9s
CI / contract (push) Successful in 1m17s
CI / build (push) Successful in 1m45s
2026-09-12 20:28:14 +07:00
Dai Ha 1db79bcaa9 Merge worker/581-completionresolver-cas-sites-0542b7-6 2026-09-12 20:28:14 +07:00
Dai Ha 4507bc5a70 Merge worker/571-attempted-outcome-5739f7-2 2026-09-12 20:28:14 +07:00
Dai Ha d7239ed23b fleetd #571: pin FleetMcp.formatReply's TIMED_OUT_UNCONFIRMED wording
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 1m42s
CORRECTION 5 on the ticket: mutating the new arm's message text to the
queued/working arm's text survived every existing test, because nothing
asserted the specific wording. This adds one test that asserts the
unconfirmed-delivery message and asserts it does NOT carry the
queued/working arm's retry invitation — the distinction #571 exists for.

No production code changes; formatReply's TIMED_OUT_UNCONFIRMED arm was
already correct.
2026-09-12 20:19:37 +07:00
Dai Ha dfeb9340b4 fleetd #581: cover completion CAS removals
CI / shell-tests (pull_request) Successful in 5s
CI / build (pull_request) Failing after 1m31s
CI / contract (pull_request) Successful in 1m35s
2026-09-12 20:19:20 +07:00
Dai Ha 1513d4f260 fleetd #562 follow-up: extract loopHealthSource factory, pin its wiring
CI / shell-tests (pull_request) Successful in 6s
CI / contract (pull_request) Successful in 59s
CI / build (pull_request) Successful in 1m43s
PR #579's inline `new FleetMcp.LoopHealthSource(poller::health, ...)` in
Fleetd.main had nothing a test could call directly. Measured: replacing
poller::health with a constant () -> RUNNING compiled clean and left all
1771 tests green (see issue #562 comment "HOLD on PR #579").

Extracts the inline construction to a package-private Fleetd.loopHealthSource
factory, the same style as the sibling capacitySource/healthCoverageSource
factories, and adds FleetdLoopHealthSourceWiringTest with three separate
assertions: the statusPoller half, the sessionReaper half, and the
reaper == null branch (still STOPPED).
2026-09-12 20:16:23 +07:00
Dai Ha 4ca7d72303 #582: assert pending message-id cleanup
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 2m32s
2026-09-12 20:15:48 +07:00
Dai Ha c0545d003d fleetd #571: make FleetApp.writeReply's inner Outcome switch exhaustive, no default
CI / shell-tests (pull_request) Successful in 8s
CI / contract (pull_request) Successful in 1m16s
CI / build (pull_request) Successful in 1m48s
Ticket comments (17126, 17127) corrected the original acceptance criterion after this
unit was already in flight: a hand-listed grep for the enum's constant names goes stale
silently the moment a new constant lands, so the compiler must be the enumeration instead.
sendOutcomeLabel (MessageService.java) and formatReply (FleetMcp.java) were already
default-free switch expressions. The one gap was writeReply's inner "status" switch, which
had `default -> "done"` — the exact value that would have lied about TIMED_OUT_UNCONFIRMED.
Remove the default and list every Outcome constant explicitly; REPLIED, COMPLETED_UNREPLIED,
QUESTION and STALE_TURN get an arm too even though the outer switch always dispatches them
first, so the inner switch stays exhaustive on its own. The outer switch (a statement, not
an expression) keeps its own default — Java does not require exhaustiveness there regardless,
and "everything not terminal is a 202" is an intentional catch-all.

Verified with the proof the ticket asked for: added a scratch 11th Outcome constant after
deleting all default arms and confirmed all three switch-expression sites (and no test file)
fail to compile without an arm for it, one at a time, then removed the scratch constant.
2026-09-12 19:28:59 +07:00
Dai Ha c1e06c9e12 fleetd #571: add TIMED_OUT_UNCONFIRMED so an ATTEMPTED delivery is not reported as never-arriving
MessageService.send's TimeoutException branch collapsed Injector.Cancellation.ATTEMPTED
(fleetd #551 — the send call was made but its outcome is unknown) into
Outcome.TIMED_OUT_QUEUED, which promises the caller the message will never arrive. On this
route agent.prompt may already have pasted and submitted the text, so a caller's natural
recovery (resend) risks a double delivery.

Add Outcome.TIMED_OUT_UNCONFIRMED and route ATTEMPTED to it. Update the three readers found
by searching for the enum's constant names (not `Outcome.`, which misses FleetMcp's
unqualified `case REPLIED ->` switches and would false-positive on ConfigRef's unrelated
Outcome record):
 - MessageService.sendOutcomeLabel: add it to the "timeout" metric label group.
 - FleetMcp.formatReply: its own case, warning against a blind retry (distinct from the
   generic "retry or poll status" message the other timeouts get).
 - FleetApp.writeReply: its own "unconfirmed" status and detail text, so it no longer falls
   through the switch's default -> "done" arm, which would have reported "the delegation
   completed" for the one case where delivery is unconfirmed.
2026-09-12 19:21:10 +07:00
18 changed files with 1320 additions and 86 deletions
+9
View File
@@ -139,6 +139,15 @@ 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` |
+19 -1
View File
@@ -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.
+143 -27
View File
@@ -24,6 +24,7 @@ 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;
@@ -80,7 +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;
@@ -501,21 +504,26 @@ public final class Fleetd {
// 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.
@@ -585,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)) {
@@ -608,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.
@@ -666,8 +660,7 @@ public final class Fleetd {
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy);
FleetMcp.LoopHealthSource loopHealth = new FleetMcp.LoopHealthSource(poller::health,
() -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health());
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)),
@@ -1042,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).
@@ -828,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]");
};
}
@@ -93,25 +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. {@link #send} reaches this outcome through {@link Injector#cancel},
* whose result tells three routes apart:
* {@link Injector.Cancellation#CANCELLED} means the message was still queued and this call
* removed it, so the target saw nothing and it will not arrive later;
* {@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 — so this route too
* establishes that the target saw nothing and it will not arrive later; but
* {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) means that call 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 hold a complete, already-submitted turn and be working on it
* right now — {@link Outcome#TIMED_OUT_WORKING}'s meaning, reported here as
* {@code TIMED_OUT_QUEUED} only because this caller never observed the pickup. Only
* {@code CANCELLED} and {@code NOT_DELIVERED} establish that the target saw nothing;
* {@code ATTEMPTED} does not.
* 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,
/**
@@ -317,20 +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 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.
* 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();
@@ -417,21 +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
* 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. {@link
* 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. 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}.
* 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);
@@ -646,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
@@ -975,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.
@@ -982,20 +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) {
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). Whatever injector.cancel() reported above — this call
// removed a still-queued Pending (CANCELLED), an earlier attempt already
// failed with a confirmed absence (NOT_DELIVERED), or an earlier attempt was
// made and its outcome is unknown (ATTEMPTED, fleetd #551 — the message may
// already have arrived in full) — the send ends with no confirmed delivery.
// 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);
@@ -670,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,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
@@ -7,6 +7,7 @@ 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;
@@ -58,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() {
@@ -299,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());
@@ -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,7 +12,11 @@ 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;
@@ -373,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 {
@@ -398,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;
}
}
@@ -601,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();
@@ -647,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();