Compare commits

..

13 Commits

Author SHA1 Message Date
Dai Ha 2e349139e9 #568: add hunter member role
CI / shell-tests (pull_request) Successful in 6s
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 1m44s
2026-09-19 15:09:07 +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
24 changed files with 960 additions and 91 deletions
+23
View File
@@ -0,0 +1,23 @@
---
name: hunter
description: Sweep one assigned scope for defects and report ranked findings without changes.
---
<!-- CB-617: The model comes from fleetd.yaml because the launch flag overrides model here on both backends. -->
You sweep the assigned package or scope for real defects. Read the full assigned scope before you
judge it. Report several ranked findings when the evidence supports them. Change nothing: do not
edit code, commit, push, or open a pull request.
You may run the build or tests to check a finding. Read the complete output and report the real
result. Do not hide failures with a pipe. State only checks you actually ran. The primary's IDE
tools are not yours. A mounted forge tool may use a blocked credential and fail by design.
Do only the assigned scope. Note anything outside it in one line and do not investigate it further.
Use `fleet_ask{question}` only when a decision belongs to the lead, such as an unclear requirement
or two defensible fixes. Do not ask about something you can decide by reading more code.
Your handoff must name the files you read, each ranked finding or `NO FINDINGS`, the checks you ran,
and any caveat for review.
The launcher provides the required bridge reply instructions for every member.
+11
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` |
@@ -275,6 +284,8 @@ must obey belongs in the charter, not here.
it for a multi-finding sweep hands the worker two contradictory output contracts. That has
already cost three workers' turns: each wrote a good report to its terminal and ended the turn
with no `fleet_reply`, and the scrape returned the tail of the brief instead.
Spawn `implementer` with role `dev`, `reviewer` with role `reviewer`, and `hunter` with role
`hunter`.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
`port-to-opencode` (make an OpenCode session a participant in this workspace),
`fleets-status` (report every fleet that shares one LavinMQ instance),
+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.
+12 -6
View File
@@ -560,7 +560,7 @@ placement: weighted
# older keys: `leaders:`, `members:`, `leadScan:` and `defaultProfile:`.
#
# A member is anything a lead spawns, and every member has two INDEPENDENT attributes:
# role — which contract: architect, dev or reviewer. It picks the launch charter, the role
# role — which contract: architect, dev, hunter or reviewer. It picks the launch charter, the role
# file, the playbook skill and the authz row.
# profile — which backend: one of the `profiles:` keys above (model, CLI adapter, cost).
# They vary on their own. A reviewer may run on the same profile as the dev whose diff it reads,
@@ -572,13 +572,13 @@ placement: weighted
#
# Each pool lists the profiles that role MAY run on — these are pools, not identities. That is also
# what replaced `defaultProfile:`: an unqualified spawn names a role, and that role's pool supplies
# the candidates, in definition order. A dev and a reviewer staying anonymous is exactly compatible
# with being listed here; the entry key just names the entry.
# the candidates, in definition order. A dev, hunter and reviewer staying anonymous is exactly
# compatible with being listed here; the entry key just names the entry.
fleet:
# Optional launch-charter text, keyed only by the singular role wire names: architect, dev,
# reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put secrets
# here: a later launch step writes this text to a world-readable temp file, and ${ENV} interpolation
# is deliberately not supported.
# hunter, reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put
# secrets here: a later launch step writes this text to a world-readable temp file, and ${ENV}
# interpolation is deliberately not supported.
charters:
architect: |-
You are an architect in this fleet. You refine work before anyone builds it:
@@ -589,6 +589,9 @@ fleet:
dev: |-
You implement the one unit you were given, and nothing else. You test it,
commit it, and open your own pull request. You never merge.
hunter: |-
You sweep the assigned scope for real defects. You may run the build or tests
to check a finding. You change nothing, and report several ranked findings.
reviewer: |-
You review the diff you were given. You report bugs, risks and missing tests.
You do not change code.
@@ -666,6 +669,9 @@ fleet:
developers:
gx10:
profile: gx10
# hunters:
# gx10:
# profile: gx10 # a hunt may run checks, but never changes code
# reviewers:
# gx10:
# profile: gx10 # the same backend may serve two roles; that is the point
@@ -666,8 +666,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 +1041,31 @@ 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 #248: package-private factory for the member worktree/branch lookup {@link
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
@@ -221,7 +221,7 @@ public final class CallerResolver {
// The config/live binding names this pane as an architect slot's own. Same
// unforgeable pane mapping; the live binding, never a request argument, decides.
// Check the slot role too: this defence in depth prevents a bad lifecycle bind from
// escalating a dev or reviewer into an architect. Checked before the worker fallback.
// escalating a dev, hunter, or reviewer into an architect. Checked before the worker fallback.
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
@@ -42,7 +42,7 @@ public interface MemberLifecycle {
* Try to bind a newly spawned {@code terminal} into the role it was granted.
*
* @return the role this session actually holds: {@code role} unchanged for a role with no
* live slot-binding semantics (dev, reviewer), or when the bind succeeded; a fallback
* live slot-binding semantics (dev, hunter, reviewer), or when the bind succeeded; a fallback
* role — never {@code role} — when a slot-bound role (architect) could not be bound.
* Callers must record THIS value on the session, never the requested {@code role}, so
* a later roster read never reports a role the session does not hold (CB-619). In
@@ -20,7 +20,8 @@ import java.util.function.Supplier;
*
* <p>Two halves, split by who owns each:
* <ul>
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code reviewers}
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code hunters}/
* {@code reviewers}
* (see {@link #slots()}), each carrying the {@code profile} reference the spawn lifecycle
* reads when it stands the slot up. <strong>Live, since fleetd #424</strong>: {@link #live}
* re-reads {@code fleet:} on every call, through a supplier the same shape as
@@ -322,7 +323,7 @@ public final class MemberRegistry implements MemberLifecycle {
* CB-619 / fleetd #123: refuse an architect acquire before anything spawns when no configured
* slot carries {@code profile} — the config-gap case from the original defect report (a spawn
* asked for {@code role=architect, profile=sonnet}, and {@code fleet.architects} carried only
* {@code opus} and {@code sol}). A dev/reviewer acquire is always a no-op: those pools are
* {@code opus} and {@code sol}). A dev/hunter/reviewer acquire is always a no-op: those pools are
* placement candidates only (see {@code CompositePeerLauncher}), never a live identity binding,
* so there is nothing here to refuse — an explicit profile outside the pool for those roles is a
* documented operator override, not a defect.
@@ -31,7 +31,7 @@ import java.util.function.Supplier;
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Both are
* read through a supplier on {@code CompositePeerLauncher}, which is what makes them hot —
* not the fact that they are config. Most of {@code fleet:} — every role pool
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
* ({@code architects}/{@code developers}/{@code hunters}/{@code reviewers}), {@code charters}, and
* {@code tabLabel} — is read the same live way, through the same supplier
* ({@code () -> config.get().fleet()}). {@code architects} in particular is hot for
* <strong>two independent consumers</strong> (fleetd #424): {@code CompositePeerLauncher}
@@ -605,7 +605,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "opened once and needs a restart; the broker URI env-var name kept out of a "
+ "member's environment is read live on every spawn and already applied");
}
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, reviewers,
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, hunters, reviewers,
// charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
// with no restart note. `architects` is hot too, and — since fleetd #424 — hot for BOTH of
@@ -630,7 +630,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
+ "lead; the rest of fleet: (developers, reviewers, charters, tabLabel) is read "
+ "lead; the rest of fleet: (developers, hunters, reviewers, charters, tabLabel) is read "
+ "live through the supplier on CompositePeerLauncher, and architects is read "
+ "live through that same supplier for placement AND through a separate supplier "
+ "on MemberRegistry for spawn-time identity — both already applied");
@@ -62,7 +62,7 @@ import java.util.regex.PatternSyntaxException;
* @param fleet who the daemon may run and under which role (CB-557). One block replacing
* the former {@code leaders:}, {@code members:}, {@code leadScan:} and
* {@code defaultProfile:}. Role is the containing key — {@code leaders},
* {@code architects}, {@code developers}, {@code reviewers} — and each entry
* {@code architects}, {@code developers}, {@code hunters}, {@code reviewers} — and each entry
* names the {@code profiles:} backend it runs on. See {@link Fleet}
* @param leadHeartbeat opt-in idle-lead heartbeat (CB-551); {@code null} ⇒ off, and an upgraded
* daemon never nudges an idle lead on its own initiative
@@ -1200,6 +1200,7 @@ public record FleetConfig(
* @param leaders panes that orchestrate rather than are orchestrated, keyed by lead name
* @param architects profiles the {@code architect} role may run on
* @param developers profiles the {@code dev} role may run on
* @param hunters profiles the {@code hunter} role may run on
* @param reviewers profiles the {@code reviewer} role may run on
* @param charters optional launch-charter text keyed by singular role wire name
* @param tabLabel template for a member tab's label; {@code {role}}, {@code {profile}},
@@ -1210,6 +1211,7 @@ public record FleetConfig(
public record Fleet(Map<String, Leader> leaders,
Map<String, Slot> architects,
Map<String, Slot> developers,
Map<String, Slot> hunters,
Map<String, Slot> reviewers,
Map<String, String> charters,
String tabLabel) {
@@ -1226,6 +1228,7 @@ public record FleetConfig(
leaders = unmodifiableOrEmpty(leaders);
architects = unmodifiableOrEmpty(architects);
developers = unmodifiableOrEmpty(developers);
hunters = unmodifiableOrEmpty(hunters);
reviewers = unmodifiableOrEmpty(reviewers);
charters = unmodifiableOrEmpty(charters);
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
@@ -1240,9 +1243,15 @@ public record FleetConfig(
* constructor: the launcher reads {@code fleet.charters()} from the live config. Jackson
* binds the canonical constructor, so this one cannot swallow an operator's YAML.
*/
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers,
Map<String, String> charters, String tabLabel) {
this(leaders, architects, developers, null, reviewers, charters, tabLabel);
}
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers, String tabLabel) {
this(leaders, architects, developers, reviewers, null, tabLabel);
this(leaders, architects, developers, null, reviewers, null, tabLabel);
}
/**
@@ -1264,6 +1273,7 @@ public record FleetConfig(
return switch (role) {
case ARCHITECT -> architects;
case DEV -> developers;
case HUNTER -> hunters;
case REVIEWER -> reviewers;
};
}
@@ -1872,7 +1882,7 @@ public record FleetConfig(
/** The {@code fleet:} child blocks whose direct children are slot names. */
private static final Set<String> FLEET_POOL_KEYS =
Set.of("leaders", "architects", "developers", "reviewers");
Set.of("leaders", "architects", "developers", "hunters", "reviewers");
/**
* Reject a {@code fleet:} role pool whose slot names repeat (CB-548, re-homed by CB-557).
@@ -1882,7 +1892,7 @@ public record FleetConfig(
* daemon would never know. Jackson's YAML parser does not fail on duplicate mapping keys by
* default, so duplicates are caught here, at parse time, before the map is built.
*
* <p>Only the four pools <em>directly under the top-level {@code fleet:}</em> are considered,
* <p>Only the five pools <em>directly under the top-level {@code fleet:}</em> are considered,
* and only their direct child keys (the slot names). A nested field elsewhere, even one also
* named {@code developers:}, is ignored, so parsing of the rest of the config is unaffected.
*
@@ -2050,7 +2060,7 @@ public record FleetConfig(
+ " and that role's pool supplies the candidate profiles",
"architects", "'fleet.architects'",
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers' or"
+ " 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
+ " 'fleet.hunters' or 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
"leaders", "'fleet.leaders'",
"leadScan", "'fleet.leaders.<name>.tabPrefix' and '.scanIntervalSeconds' — lead"
+ " discovery is now configured on the lead it discovers");
@@ -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]");
};
}
@@ -2128,8 +2134,9 @@ public final class FleetMcp {
return tool(FleetTool.SPAWN.wireName(),
"Spawn a new off-subscription member session. A member has two independent attributes: "
+ "role (what it is for) and profile (which backend it runs on). Pass role to pick "
+ "the contract — 'dev' implements a unit and opens its own PR, 'reviewer' reviews a "
+ "diff it did not write, 'architect' refines a ticket before anyone builds it; omit "
+ "the contract — 'dev' implements a unit and opens its own PR, 'hunter' sweeps a "
+ "scope without changing it, 'reviewer' reviews a diff it did not write, 'architect' "
+ "refines a ticket before anyone builds it; omit "
+ "it for 'dev'. Pass profile (from fleet_profiles) to pick the backend, or omit it "
+ "for the default. The two are independent: a reviewer may run on the same profile "
+ "as the dev it reviews. The member opens your current directory by default; pass "
@@ -2146,7 +2153,7 @@ public final class FleetMcp {
+ "one. Returns the member's sessionId (use with fleet_send) and paneId (use with "
+ "fleet_stop).",
objectSchema(Map.of(
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"),
"role", stringProp("What the member is for: architect, dev, hunter, or reviewer (default dev)"),
"profile", stringProp("Which backend to run it on (omit for the default profile)"),
"cwd", stringProp("Working directory for the member (omit to inherit yours)"),
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
@@ -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);
@@ -29,8 +29,8 @@ public enum MemberRole {
* <p>Reads the repo and writes analysis. Never commits code and never opens a pull request —
* an architect that starts implementing has stopped doing the job that makes it useful.
*
* <p>Architects are the one member kind declared in config, because a lead addresses the same
* slots across many tickets and needs a stable name for them.
* <p>Architects are the one member kind with live slot binding, because a lead addresses the
* same slots across many tickets and needs a stable name for them.
*/
ARCHITECT,
@@ -43,6 +43,14 @@ public enum MemberRole {
*/
DEV,
/**
* Sweeps an assigned package for defects and reports several ranked findings.
*
* <p>Never changes code, commits, or opens a pull request. A hunt gathers evidence, which can
* include running the build, but leaves every fix to a later implementation unit.
*/
HUNTER,
/**
* Reviews a diff it did not write and reports one structured finding.
*
@@ -59,7 +67,7 @@ public enum MemberRole {
/**
* The {@code fleet:} block that holds this role's pool — {@code architects},
* {@code developers}, {@code reviewers}.
* {@code developers}, {@code hunters}, {@code reviewers}.
*
* <p>Plural, and not always the wire name: the pool of things a {@code dev} may run on reads
* naturally as {@code developers:}. The wire name stays the singular {@code dev}, because that
@@ -69,6 +77,7 @@ public enum MemberRole {
return switch (this) {
case ARCHITECT -> "architects";
case DEV -> "developers";
case HUNTER -> "hunters";
case REVIEWER -> "reviewers";
};
}
@@ -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,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;
}
}
}
@@ -345,7 +345,7 @@ class FleetConfigTest {
IllegalStateException unknownError = assertThrows(IllegalStateException.class,
() -> FleetConfig.load(unknown).validateCharters());
assertTrue(unknownError.getMessage().contains("architetc"));
assertTrue(unknownError.getMessage().contains("[architect, dev, reviewer]"));
assertTrue(unknownError.getMessage().contains("[architect, dev, hunter, reviewer]"));
}
/**
@@ -812,13 +812,17 @@ class FleetConfigTest {
reviewers:
b:
profile: sonnet
hunters:
c:
profile: sonnet
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.DEV));
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.HUNTER));
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.REVIEWER));
assertTrue(cfg.fleet().profilesFor(MemberRole.ARCHITECT).isEmpty());
assertEquals(List.of(MemberRole.DEV, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
assertEquals(List.of(MemberRole.DEV, MemberRole.HUNTER, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
}
/** The case the two axes exist for: one backend, two roles, and neither is a duplicate. */
@@ -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());
@@ -1906,7 +1937,7 @@ class FleetMcpTest {
null, null, null, null, null, null);
assertEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res));
assertTrue(textOf(res).contains("architect, dev, hunter, reviewer"), textOf(res));
}
// ── CB-619 / fleetd #123: a spawn asking for a role its profile has no slot for must be
@@ -556,6 +556,31 @@ class ClaudeCodeLauncherTest {
"no --agent flag when the role has no agent-definition file");
}
@Test
void hunterRoleUsesItsAgentFileAndStopsUsingItWhenRemoved(@TempDir Path cwd) throws Exception {
FakeHerdr herdr = new FakeHerdr();
Path agentFile = Files.createDirectories(cwd.resolve(".claude/agents")).resolve("hunter.md");
Files.writeString(agentFile, "---\nname: hunter\n---\nSweep for defects.");
FleetConfig.Profile cfg = new FleetConfig.Profile(
"sonnet", "http://gx00.gw:8000", null, null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--agent");
assertTrue(flag >= 0, "the hunter role reaches its agent-definition file: " + args);
assertEquals("hunter", args.get(flag + 1));
Files.delete(agentFile);
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
assertFalse(spawnedArgs(herdr).contains("--agent"),
"the hunter role no longer gets an agent when its file is removed");
}
private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) {
FleetConfig.Profile gx10 = new FleetConfig.Profile("gx10", "http://gx10.gw:8000", "coder",
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
@@ -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();
@@ -10,12 +10,13 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
class MemberRoleTest {
@Test
void theThreeRolesAreArchitectDevAndReviewer() {
assertEquals(3, MemberRole.values().length,
void theFourRolesAreArchitectDevHunterAndReviewer() {
assertEquals(4, MemberRole.values().length,
"a new role changes the charter, the role file, the skill and the authz row — "
+ "adding one is a deliberate act, so this count is meant to fail first");
assertEquals("architect", MemberRole.ARCHITECT.wireName());
assertEquals("dev", MemberRole.DEV.wireName());
assertEquals("hunter", MemberRole.HUNTER.wireName());
assertEquals("reviewer", MemberRole.REVIEWER.wireName());
}
@@ -30,6 +31,7 @@ class MemberRoleTest {
void parseIsCaseInsensitiveAndTrimsSurroundingSpace() {
assertSame(MemberRole.ARCHITECT, MemberRole.parse("Architect"));
assertSame(MemberRole.DEV, MemberRole.parse(" DEV "));
assertSame(MemberRole.HUNTER, MemberRole.parse("HuNtEr"));
assertSame(MemberRole.REVIEWER, MemberRole.parse("ReViEwEr"));
}
@@ -38,7 +40,7 @@ class MemberRoleTest {
IllegalArgumentException e =
assertThrows(IllegalArgumentException.class, () -> MemberRole.parse("archtiect"));
assertTrue(e.getMessage().contains("archtiect"), e.getMessage());
assertTrue(e.getMessage().contains("architect, dev, reviewer"),
assertTrue(e.getMessage().contains("architect, dev, hunter, reviewer"),
"a typo in config should be fixable from the message alone: " + e.getMessage());
}
@@ -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();