Compare commits
29 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 856dfc6318 | |||
| 81c1d8e91c | |||
| f5c6a0e4fc | |||
| a7aee5b982 | |||
| bf895616a5 | |||
| b874afb0af | |||
| 6eb34a654f | |||
| 7084d99b89 | |||
| ae7845c375 | |||
| 61115f6f61 | |||
| 4b9ebda1b3 | |||
| 1e68d7ee39 | |||
| 6cccd458d4 | |||
| 42820fbe75 | |||
| d91ff886da | |||
| 386e760a5c | |||
| 2e349139e9 | |||
| a639969a9a | |||
| 17c3a69c57 | |||
| 49a5875586 | |||
| 634d33b50b | |||
| 1db79bcaa9 | |||
| 4507bc5a70 | |||
| d7239ed23b | |||
| dfeb9340b4 | |||
| 1513d4f260 | |||
| 4ca7d72303 | |||
| c0545d003d | |||
| c1e06c9e12 |
@@ -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.
|
||||
@@ -139,6 +139,16 @@ 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. Architects first form independent positions, then
|
||||
compare them. If they still disagree after that comparison, they return both positions and their
|
||||
checked evidence; the lead decides. Go to the operator only for an action the fleet has no
|
||||
authority to take, such as spending money, granting access, or making a promise 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 +285,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),
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
@@ -218,23 +221,23 @@ public final class Fleetd {
|
||||
// there is no 2-arg overload left for any lambda to silently bind to instead), but so a
|
||||
// test can call the exact same object this line builds, instead of asserting a copy of its
|
||||
// shape (round 3's lesson).
|
||||
ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
|
||||
// fleetd #589 Group 1: extracted to forwardingExhaustionSink(...) below (see that method's
|
||||
// javadoc) so a dedicated test can prove this factory keeps reading the reference live,
|
||||
// rather than a rebuilt copy of its shape.
|
||||
ExhaustionSink forwardingExhaustionSink = forwardingExhaustionSink(exhaustionSinkRef);
|
||||
// The claude-code adapter is the always-present default; keep it even with no profiles (so a
|
||||
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
|
||||
// unless opencode is the only kind configured.
|
||||
// fleetd #589 Group 2: extracted to claudeCodeLauncher(...)/openCodeLauncher(...) below (see
|
||||
// those methods' javadoc) so a dedicated test can prove the CB-596 memberCredentials policy
|
||||
// supplier is actually wired to each adapter, not silently replaced with `() -> null`.
|
||||
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
|
||||
adapters.add(new ClaudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
|
||||
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials(), null, config::get));
|
||||
adapters.add(claudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
|
||||
claudeProfiles, cfg, config));
|
||||
}
|
||||
if (!opencodeProfiles.isEmpty()) {
|
||||
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
|
||||
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials(), config::get, forwardingExhaustionSink));
|
||||
adapters.add(openCodeLauncher(router.memberAgents(), router.memberSpaces(),
|
||||
opencodeProfiles, cfg, config, forwardingExhaustionSink));
|
||||
}
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||
@@ -408,12 +411,12 @@ public final class Fleetd {
|
||||
// LiveExhaustedPatterns's class doc for why this replaces the old compiled-once-at-startup
|
||||
// map. A profile with no exhaustedPattern simply returns null here, so its workers keep
|
||||
// today's completion-fallback behaviour unchanged.
|
||||
LiveExhaustedPatterns liveExhaustedPatterns = new LiveExhaustedPatterns(() -> config.get().profiles());
|
||||
ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> liveExhaustedPatterns.patternFor(session.profile()))
|
||||
.orElse(null);
|
||||
// fleetd #589 Group 1: both extracted to liveExhaustedPatterns(...)/
|
||||
// exhaustedPatternLookup(...) below (see those methods' javadoc) — this is the worst
|
||||
// consequence in the whole #589 sweep: silently losing either wiring means a genuine
|
||||
// usage-limit refusal is handed back as a real completion instead of BACKEND_EXHAUSTED.
|
||||
LiveExhaustedPatterns liveExhaustedPatterns = liveExhaustedPatterns(config);
|
||||
ExhaustedPatternLookup exhaustedPatterns = exhaustedPatternLookup(sessions::roster, liveExhaustedPatterns);
|
||||
// The startup coverage line still reports the boot-time snapshot only — it is printed once,
|
||||
// here, and a reload no longer needs to change what it said; exhaustionDetectionArmed (via
|
||||
// liveExhaustedPatterns.armed, wired into quarantineSource below) is what stays live.
|
||||
@@ -461,11 +464,13 @@ public final class Fleetd {
|
||||
// method's javadoc for the full fleetd #175/#234/#446 history this used to carry inline —
|
||||
// so a dedicated test can drive the exact ExhaustionSink main() builds, not a hand-rebuilt
|
||||
// copy of its shape.
|
||||
ExhaustionSink exhaustionSink = exhaustionSink(sessions, config, quarantine,
|
||||
quarantineReasonByCredential, cfg);
|
||||
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
|
||||
// now that `sessions` exists to resolve target -> session -> profile.
|
||||
exhaustionSinkRef.set(exhaustionSink);
|
||||
// fleetd #589 Group 1: both statements (build + set) folded into publishExhaustionSink(...)
|
||||
// below (see that method's javadoc), so a test can prove the reference is actually
|
||||
// repointed at the real sink, not silently left at ExhaustionSink.none().
|
||||
ExhaustionSink exhaustionSink = publishExhaustionSink(exhaustionSinkRef, sessions, config,
|
||||
quarantine, quarantineReasonByCredential, cfg);
|
||||
// fleetd #201 Unit 5: the production BackendErrorSink needs `pushLoop` (built further below,
|
||||
// after `sessions`) to tell a lead about an incident or an unmapped target — the same
|
||||
// construction-order cycle `exhaustionSinkRef` breaks above, broken the same way: a mutable
|
||||
@@ -501,21 +506,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 +595,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 +620,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 +662,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)),
|
||||
@@ -1010,6 +1005,120 @@ public final class Fleetd {
|
||||
cfg.profiles()::keySet, System::nanoTime);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: the forwarding {@link ExhaustionSink} handed to the adapters built
|
||||
* before {@code sessions} exists (see the {@code exhaustionSinkRef}/{@code
|
||||
* forwardingExhaustionSink} locals in {@code main}, just above {@link #capacitySource}'s call
|
||||
* site). Before this ticket, {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)} was
|
||||
* built inline — nothing a test could call directly, so a mutation swapping the supplier for a
|
||||
* hardcoded {@code () -> ExhaustionSink.none()} compiled clean and left the suite green: the
|
||||
* forwarder would silently stop reading the reference at all, and {@link
|
||||
* #publishExhaustionSink} repointing that reference later would have no effect.
|
||||
*
|
||||
* <p>Extracted the same way {@link #capacitySource}/{@link #loopHealthSource} were, so {@code
|
||||
* FleetdExhaustionSinkForwardingWiringTest} can call this factory directly with a real {@link
|
||||
* AtomicReference}, mutate the reference AFTER the forwarder is built, and prove the forwarder
|
||||
* still reads it live rather than a fixed target captured at construction time.
|
||||
*/
|
||||
static ExhaustionSink forwardingExhaustionSink(AtomicReference<ExhaustionSink> exhaustionSinkRef) {
|
||||
return ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: publish the real {@link ExhaustionSink} — built the same way {@link
|
||||
* #exhaustionSink} always was — into the forwarding reference {@link #forwardingExhaustionSink}
|
||||
* built above, replacing {@code main}'s previously untested two-statement sequence ({@code
|
||||
* ExhaustionSink exhaustionSink = exhaustionSink(...); exhaustionSinkRef.set(exhaustionSink);}).
|
||||
* Before this ticket, nothing proved the {@code .set(...)} call actually received the real sink
|
||||
* rather than a hardcoded {@code ExhaustionSink.none()} — the whole point of {@code
|
||||
* exhaustionSinkRef} existing (fleetd #175) is that {@link
|
||||
* dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch check, built before {@code sessions}
|
||||
* exists, keeps working once this line runs; silently keeping the reference at {@code none()}
|
||||
* would mean that check permanently does nothing, with the full suite still green because no
|
||||
* existing test drives this exact call site.
|
||||
*
|
||||
* <p>Returns the built sink so {@code main} can still pass it to {@link CompletionResolver}'s
|
||||
* constructor at the same call site it already does, without building it twice.
|
||||
*/
|
||||
static ExhaustionSink publishExhaustionSink(AtomicReference<ExhaustionSink> exhaustionSinkRef,
|
||||
SessionManager sessions, ConfigRef config, BackendQuarantine quarantine,
|
||||
Map<String, String> quarantineReasonByCredential, FleetConfig cfg) {
|
||||
ExhaustionSink sink = exhaustionSink(sessions, config, quarantine, quarantineReasonByCredential, cfg);
|
||||
exhaustionSinkRef.set(sink);
|
||||
return sink;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: the live {@code exhaustedPattern} source (fleetd #446) {@link
|
||||
* CompletionResolver} enforces on, extracted out of {@code main} for the same reason {@link
|
||||
* #capacitySource} was. Before this ticket {@code new LiveExhaustedPatterns(() ->
|
||||
* config.get().profiles())} was built inline; replacing the supplier with a hardcoded {@code ()
|
||||
* -> Map.of()} compiled clean and left the suite green, meaning every profile's {@code
|
||||
* exhaustedPattern} would silently stop being recognised and a genuine usage-limit refusal
|
||||
* would be handed back as a real completion instead of {@code BACKEND_EXHAUSTED}.
|
||||
*/
|
||||
static LiveExhaustedPatterns liveExhaustedPatterns(ConfigRef config) {
|
||||
return new LiveExhaustedPatterns(() -> config.get().profiles());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: the {@link ExhaustedPatternLookup} {@link CompletionResolver} enforces
|
||||
* on, resolving a herdr {@code target} to its session's profile and then to that profile's live
|
||||
* {@link LiveExhaustedPatterns#patternFor}. Extracted out of {@code main} the same way {@link
|
||||
* #worktreeBranchLookup} was — same {@code Supplier<List<MemberSession>>} roster shape, same
|
||||
* reason: before this ticket the lambda was built inline, and replacing it with {@code target ->
|
||||
* null} (the exact shape of {@link ExhaustedPatternLookup#none()}) compiled clean and left the
|
||||
* suite green. This is the worst consequence in the whole #589 sweep (see the ticket): a
|
||||
* genuine usage-limit refusal would stop being classified as {@code BACKEND_EXHAUSTED} and
|
||||
* would be handed back to a waiting {@code fleet_send} as if it were real completed work.
|
||||
*
|
||||
* @param roster the live member roster, normally {@code sessions::roster}
|
||||
*/
|
||||
static ExhaustedPatternLookup exhaustedPatternLookup(Supplier<List<MemberSession>> roster,
|
||||
LiveExhaustedPatterns liveExhaustedPatterns) {
|
||||
return target -> roster.get().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> liveExhaustedPatterns.patternFor(session.profile()))
|
||||
.orElse(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 2: the production {@link ClaudeCodeLauncher} adapter, extracted out of
|
||||
* {@code main} the same way {@link #capacitySource} was. Before this ticket the constructor
|
||||
* call (11 arguments, including the CB-596 {@code memberCredentials} policy supplier) was built
|
||||
* inline; replacing the {@code () -> config.get().memberCredentials()} argument with {@code ()
|
||||
* -> null} compiled clean and left the suite green — {@code memberCredentials} is not {@code
|
||||
* null} itself (a lambda is never {@code null}), so {@link
|
||||
* dev.ltms.fleet.member.HerdrPeerLauncher#applyMemberCredentialPolicy} sees {@code
|
||||
* memberCredentials.get() == null} and silently shadows nothing, reopening the exact CB-592
|
||||
* exposure gap CB-596's policy closed. {@code FleetdClaudeCodeLauncherCredentialWiringTest}
|
||||
* calls this factory with a real {@link ConfigRef} carrying a {@code memberCredentials:} block
|
||||
* and proves a known-but-not-allowed name is actually shadowed on {@code spawn()}.
|
||||
*/
|
||||
static ClaudeCodeLauncher claudeCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||
SubscriptionGuard guard, Map<String, FleetConfig.Profile> claudeProfiles, FleetConfig cfg,
|
||||
ConfigRef config) {
|
||||
return new ClaudeCodeLauncher(agents, spaces, guard, claudeProfiles, cfg.effectiveDefaultProfile(),
|
||||
System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(), () -> config.get().memberCredentials(), null, config::get);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 2: the production {@link OpenCodeLauncher} adapter, the {@code opencode}
|
||||
* counterpart to {@link #claudeCodeLauncher} above and extracted for the identical reason: the
|
||||
* same {@code () -> config.get().memberCredentials()} argument, reopening the same CB-592
|
||||
* exposure gap if silently replaced with {@code () -> null}.
|
||||
*/
|
||||
static OpenCodeLauncher openCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||
Map<String, FleetConfig.Profile> opencodeProfiles, FleetConfig cfg, ConfigRef config,
|
||||
ExhaustionSink forwardingExhaustionSink) {
|
||||
return new OpenCodeLauncher(agents, spaces, opencodeProfiles, cfg.effectiveDefaultProfile(),
|
||||
System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(), () -> config.get().memberCredentials(), config::get,
|
||||
forwardingExhaustionSink);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #426: package-private factory for {@code fleet_list}'s {@code healthCoverage} source,
|
||||
* extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
|
||||
@@ -1042,6 +1151,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).
|
||||
|
||||
@@ -221,7 +221,8 @@ 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,8 @@ 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.
|
||||
*
|
||||
@@ -2049,8 +2059,9 @@ public record FleetConfig(
|
||||
"defaultProfile", "a role pool under 'fleet:' — an unqualified spawn now names a role,"
|
||||
+ " 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",
|
||||
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers',"
|
||||
+ " '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");
|
||||
|
||||
@@ -10,6 +10,8 @@ import java.io.IOException;
|
||||
import java.io.UncheckedIOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Collections;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
@@ -25,9 +27,11 @@ import java.util.function.Supplier;
|
||||
* every gate ({@link #confirm}'s own checks) has passed — a deferred, single-shot continuation
|
||||
* clears the lead's own pane and bootstraps a fresh session against that file.
|
||||
*
|
||||
* <p>This is the executor only. Nothing in this ticket wires an MCP tool onto {@link #open}/
|
||||
* {@link #confirm}/{@link #cancel} — that is a separate, later unit; until it lands, nothing calls
|
||||
* this class at all.
|
||||
* <p>This is the executor behind the {@code fleet_handover} MCP tool ({@code
|
||||
* dev.ltms.fleet.mcp.FleetMcp#handover}), which drives {@link #open}, {@link #confirm}, {@link
|
||||
* #cancel}, and {@link #status} from a tool call — wired in fleetd #480 Unit C. <strong>An earlier
|
||||
* version of this paragraph said nothing called this class at all; that stopped being true once
|
||||
* that unit landed, and this correction exists so the javadoc does not go on claiming it.</strong>
|
||||
*
|
||||
* <p><strong>{@code confirm()} cannot roll inline — a fleetd #480 correction.</strong> The first
|
||||
* version of this class called {@code agents.send(lead, "/clear")} directly from inside {@code
|
||||
@@ -181,6 +185,91 @@ public final class LeadRollover {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* How many tokens {@link #outcomes} remembers before it starts evicting the oldest — bounded
|
||||
* so a long-running daemon never grows this map without limit. Chosen generously rather than
|
||||
* tightly: production rolls are rare (this class's own ticket found exactly ONE completed roll
|
||||
* ever logged on this host), and each entry is a handful of short strings, so even a full cap
|
||||
* costs a few tens of kilobytes — nowhere near a reason to make it configurable. 200 entries
|
||||
* comfortably outlasts any operator's own memory of "did that roll I asked for actually
|
||||
* happen", which is the whole reason {@link #status} exists.
|
||||
*
|
||||
* <p><strong>This cap counts {@link RollState#IN_PROGRESS} entries exactly the same as
|
||||
* finished ones.</strong> There is only the one bounded map: {@link #confirm} writes an {@link
|
||||
* RollState#IN_PROGRESS} entry into {@link #outcomes} at hand-off, and the deferred
|
||||
* continuation later overwrites that SAME key with a terminal state — it never inserts a
|
||||
* second entry. An approved roll therefore occupies one slot in this map for its entire
|
||||
* lifetime, from the moment {@link #confirm} hands off, not only once it finishes; a
|
||||
* confirmed-but-not-yet-finished roll counts against the cap exactly like a finished one. The
|
||||
* alternative (a separate, uncapped in-flight map) would let a burst of confirmed-but-stuck
|
||||
* rolls grow without bound — the exact failure this cap exists to prevent — so it was rejected.
|
||||
*/
|
||||
static final int OUTCOME_HISTORY_CAP = 200;
|
||||
|
||||
/**
|
||||
* What is known about one token, right now — the answer {@link #status} gives. Distinguishes
|
||||
* three terminal outcomes an approved roll can finish with, one in-flight outcome for a roll
|
||||
* that has been approved but has not finished yet, and two answers for a token that names no
|
||||
* active work at all: still pending confirmation, or nothing known about this token at all.
|
||||
*/
|
||||
public enum RollState {
|
||||
/**
|
||||
* {@code token} is still open: either {@link #open} was called and {@link #confirm} has not
|
||||
* been (or not successfully) yet, or a {@link #confirm} call failed one of its gate checks
|
||||
* and left the token pending for a retry — see {@link #confirm}'s javadoc ("token stays
|
||||
* pending"). Indistinguishable from a genuinely fresh request; a caller wanting to know
|
||||
* WHICH gate most recently refused should read the {@link RollDecision} that {@link
|
||||
* #confirm} itself returned, not this status. <strong>Never the state of an APPROVED
|
||||
* roll</strong> — see {@link #IN_PROGRESS}, which {@link #confirm} records at the moment it
|
||||
* hands off, before this token is even removed from the pending set.
|
||||
*/
|
||||
PENDING,
|
||||
/**
|
||||
* {@link #confirm} approved this roll and handed it to the deferred continuation, which has
|
||||
* not finished yet. Recorded by {@link #confirm} itself, at hand-off — <strong>before</strong>
|
||||
* {@code token} is removed from the pending set — so there is never a gap in which {@link
|
||||
* #status} could wrongly answer {@link #UNKNOWN} ("nothing was ever requested") for a roll
|
||||
* that is, in fact, actively running. This is not sticky: the deferred continuation
|
||||
* overwrites this same entry with a terminal state ({@link #ROLLED}, {@link
|
||||
* #TURN_NEVER_SETTLED}, or {@link #CLEAR_NEVER_SETTLED}) once it finishes.
|
||||
*/
|
||||
IN_PROGRESS,
|
||||
/**
|
||||
* {@link #confirm} was approved and the deferred continuation completed the entire roll:
|
||||
* the calling lead's turn settled, {@code /clear} was sent and settled, and {@code
|
||||
* bootstrapText} was sent.
|
||||
*/
|
||||
ROLLED,
|
||||
/**
|
||||
* {@link #confirm} was approved, but the calling lead's own turn never reached a boundary
|
||||
* (IDLE or DONE) within {@code turnSettleSeconds} — no {@code /clear} was ever sent, at
|
||||
* all. This is the branch the fleetd #480 correction exists to make safe, and the one this
|
||||
* status exists to make VISIBLE: before this, a lead that hit this case had no way to find
|
||||
* out, and would carry on believing it was about to be replaced. See this class's javadoc.
|
||||
*/
|
||||
TURN_NEVER_SETTLED,
|
||||
/**
|
||||
* {@link #confirm} was approved and {@code /clear} was sent, but the pane never re-settled
|
||||
* within {@code clearSettleSeconds} — {@code bootstrapText} was never sent.
|
||||
*/
|
||||
CLEAR_NEVER_SETTLED,
|
||||
/**
|
||||
* {@code token} names nothing this instance currently knows about: never issued by {@link
|
||||
* #open}, dropped by {@link #cancel}, or aged out of {@link #outcomes}'s bounded history.
|
||||
* These three causes are not distinguished — all of them mean "there is nothing to tell
|
||||
* you", which is the entire content of a clean answer here.
|
||||
*/
|
||||
UNKNOWN
|
||||
}
|
||||
|
||||
/**
|
||||
* The answer {@link #status} gives for one token: a {@link RollState} and a human-readable
|
||||
* {@code detail}. For {@link RollState#TURN_NEVER_SETTLED}, {@code detail} names {@code
|
||||
* turnSettleSeconds} and its configured value explicitly, so a reader who sees this knows what
|
||||
* to raise.
|
||||
*/
|
||||
public record RollStatus(RollState state, String detail) {}
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Supplier<FleetConfig.LeadRollover> configSupplier;
|
||||
/**
|
||||
@@ -201,6 +290,23 @@ public final class LeadRollover {
|
||||
*/
|
||||
private final Consumer<Runnable> continuationRunner;
|
||||
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Finished tokens → what actually happened, for {@link #status}. Bounded by {@link
|
||||
* #OUTCOME_HISTORY_CAP}, oldest evicted first ({@code removeEldestEntry} on an insertion-order
|
||||
* {@link LinkedHashMap}). Wrapped in {@link Collections#synchronizedMap} because entries are
|
||||
* written from whatever thread {@code continuationRunner} runs the roll on (a fresh virtual
|
||||
* thread in production, the calling test thread under {@code Runnable::run}) and read from
|
||||
* whatever thread calls {@link #status} (the MCP handler thread) — a plain {@code
|
||||
* LinkedHashMap} is not safe for that, and {@code removeEldestEntry} additionally requires
|
||||
* external synchronization even for a thread-safe map that merely wraps it.
|
||||
*/
|
||||
private final Map<String, RollStatus> outcomes = Collections.synchronizedMap(
|
||||
new LinkedHashMap<>(16, 0.75f, false) {
|
||||
@Override
|
||||
protected boolean removeEldestEntry(Map.Entry<String, RollStatus> eldest) {
|
||||
return size() > OUTCOME_HISTORY_CAP;
|
||||
}
|
||||
});
|
||||
|
||||
/** Production constructor — wall clock, real sleep between settle polls, a real virtual thread. */
|
||||
public LeadRollover(AgentControl agents, Supplier<FleetConfig.LeadRollover> configSupplier,
|
||||
@@ -366,6 +472,15 @@ public final class LeadRollover {
|
||||
return docCheck;
|
||||
}
|
||||
|
||||
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
|
||||
// OUTCOME_HISTORY_CAP's javadoc. This ordering means `token` is written into `outcomes`
|
||||
// while it is STILL present in `pending`; status() checks `outcomes` first (see that
|
||||
// method), so it reports IN_PROGRESS immediately, not the brief-but-real gap a
|
||||
// remove-then-put ordering would leave in which the token is in neither map.
|
||||
outcomes.put(token, new RollStatus(RollState.IN_PROGRESS,
|
||||
"confirm() approved this roll and handed it to the deferred continuation; it has "
|
||||
+ "not finished yet — still waiting for the calling turn to settle, for "
|
||||
+ "/clear to be sent and settle, or for bootstrapText to be sent"));
|
||||
pending.remove(token);
|
||||
log.info("lead-rollover: confirmed token={} lead={} — roll scheduled once the calling turn ends",
|
||||
token, callerTerminal);
|
||||
@@ -393,6 +508,11 @@ public final class LeadRollover {
|
||||
+ "turn is still live and clearing it now would destroy live context "
|
||||
+ "(token={}, configured={}s elapsed={}ms)",
|
||||
lead, p.token(), cfg.turnSettleSeconds(), turnResult.elapsedMillis());
|
||||
outcomes.put(p.token(), new RollStatus(RollState.TURN_NEVER_SETTLED,
|
||||
"the calling lead's own turn never reached a boundary (IDLE or DONE) within "
|
||||
+ "turnSettleSeconds=" + cfg.turnSettleSeconds() + "s (measured elapsed="
|
||||
+ turnResult.elapsedMillis() + "ms) — no /clear was ever sent. If this "
|
||||
+ "keeps happening, raise turnSettleSeconds in fleetd.yaml"));
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -411,11 +531,17 @@ public final class LeadRollover {
|
||||
+ "elapsed={}ms nudges={})",
|
||||
lead, p.token(), cfg.clearSettleSeconds(), clearResult.elapsedMillis(),
|
||||
clearResult.nudges());
|
||||
outcomes.put(p.token(), new RollStatus(RollState.CLEAR_NEVER_SETTLED,
|
||||
"/clear was sent, but the pane never re-settled within clearSettleSeconds="
|
||||
+ cfg.clearSettleSeconds() + "s (measured elapsed=" + clearResult.elapsedMillis()
|
||||
+ "ms, nudges=" + clearResult.nudges() + ") — bootstrapText was never sent"));
|
||||
return;
|
||||
}
|
||||
agents.send(lead, cfg.bootstrapTextFor(p.handoverPath()));
|
||||
long rollElapsedMillis = nowMillis.getAsLong() - rollStartMillis;
|
||||
log.info("lead-rollover: rolled token={} lead={} elapsedMs={}", p.token(), lead, rollElapsedMillis);
|
||||
outcomes.put(p.token(), new RollStatus(RollState.ROLLED,
|
||||
"rolled successfully in " + rollElapsedMillis + "ms"));
|
||||
}
|
||||
|
||||
/** Drop a pending request without rolling. @return whether a pending request existed for {@code token} */
|
||||
@@ -423,6 +549,47 @@ public final class LeadRollover {
|
||||
return pending.remove(token) != null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Read-only: what is currently known about {@code token}. <strong>Never sends anything, never
|
||||
* schedules, cancels, or retries a roll</strong> — a caller may poll this as often as it likes
|
||||
* with no side effect at all, which is exactly why it exists: every failure past {@link
|
||||
* #confirm} used to be a {@code log.warn} a lead can never read (see this class's javadoc), and
|
||||
* this is the only route back.
|
||||
*
|
||||
* @param token the token {@link #open} returned; {@code null} or blank is a clean {@link
|
||||
* RollState#UNKNOWN}, never a {@link NullPointerException} — {@link #pending} is a
|
||||
* {@link ConcurrentHashMap}, which throws on a {@code null} key lookup, so this
|
||||
* short-circuits before ever reaching it
|
||||
* @return {@link RollState#IN_PROGRESS} for an approved roll whose continuation has not
|
||||
* finished yet, or a terminal state once it has (both read from {@link #outcomes} —
|
||||
* checked FIRST, see below); {@link RollState#PENDING} while {@code token} is still
|
||||
* open and has not yet been approved (including one left pending by a {@link #confirm}
|
||||
* gate refusal — see that method's javadoc); or {@link RollState#UNKNOWN} for a token
|
||||
* never issued, cancelled, or aged out of the bounded history
|
||||
*/
|
||||
public RollStatus status(String token) {
|
||||
if (token == null || token.isBlank()) {
|
||||
return new RollStatus(RollState.UNKNOWN, "no token given");
|
||||
}
|
||||
// `outcomes` is checked BEFORE `pending`, deliberately: `confirm` writes an IN_PROGRESS
|
||||
// entry into `outcomes` before it removes `token` from `pending` (see `confirm`'s own
|
||||
// comment at that call site), so for the brief window where a token is present in BOTH
|
||||
// maps, this order reports the more accurate answer (IN_PROGRESS, already approved) rather
|
||||
// than the stale one (PENDING, not yet approved) a pending-first check would give.
|
||||
RollStatus recorded = outcomes.get(token);
|
||||
if (recorded != null) {
|
||||
return recorded;
|
||||
}
|
||||
if (pending.containsKey(token)) {
|
||||
return new RollStatus(RollState.PENDING, "open() has been called for this token and "
|
||||
+ "it has not yet been confirmed — or a confirm() gate check failed and left it "
|
||||
+ "pending, so the same token may be retried once the problem is fixed");
|
||||
}
|
||||
return new RollStatus(RollState.UNKNOWN, "token names no pending or finished rollover "
|
||||
+ "request known to this instance — never issued, cancelled, or aged out of the "
|
||||
+ "bounded history (cap=" + OUTCOME_HISTORY_CAP + ")");
|
||||
}
|
||||
|
||||
/**
|
||||
* The three handover-file checks, in order: exists, not empty, fresh (modified after
|
||||
* {@link #open}'s timestamp and not older than {@code maxDocAgeSeconds}). Stats {@code
|
||||
|
||||
@@ -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]");
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1242,14 +1248,16 @@ public final class FleetMcp {
|
||||
Map<String, Object> args) {
|
||||
String action = str(args, "action");
|
||||
if (isBlank(action)) {
|
||||
return error("action is required: \"open\", \"confirm\" or \"cancel\"");
|
||||
return error("action is required: \"open\", \"confirm\", \"cancel\" or \"status\"");
|
||||
}
|
||||
return switch (action) {
|
||||
case "open" -> handoverOpen(leadRollover, callerTerminal, str(args, "reason"));
|
||||
case "confirm" -> handoverConfirm(leadRollover, callerTerminal, str(args, "token"),
|
||||
truthy(args, "operatorConfirmed"));
|
||||
case "cancel" -> handoverCancel(leadRollover, str(args, "token"));
|
||||
default -> error("unknown action \"" + action + "\" — must be \"open\", \"confirm\" or \"cancel\"");
|
||||
case "status" -> handoverStatus(leadRollover, str(args, "token"));
|
||||
default -> error("unknown action \"" + action
|
||||
+ "\" — must be \"open\", \"confirm\", \"cancel\" or \"status\"");
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1319,6 +1327,29 @@ public final class FleetMcp {
|
||||
return text(json(m));
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code action: "status"}. Read-only — see {@link LeadRollover#status}: never schedules,
|
||||
* cancels, or retries anything, and it is the only way for a lead to find out what happened to
|
||||
* a token past {@code confirm()}, since every outcome after that point is otherwise logged only
|
||||
* (see {@link LeadRollover}'s class javadoc).
|
||||
*/
|
||||
private static McpSchema.CallToolResult handoverStatus(LeadRollover leadRollover, String token) {
|
||||
if (leadRollover == null) {
|
||||
Map<String, Object> m = new LinkedHashMap<>();
|
||||
m.put("state", "NOT_CONFIGURED");
|
||||
m.put("detail", "leadRollover: is not configured");
|
||||
return text(json(m));
|
||||
}
|
||||
if (isBlank(token)) {
|
||||
return error("token is required for action \"status\"");
|
||||
}
|
||||
LeadRollover.RollStatus s = leadRollover.status(token);
|
||||
Map<String, Object> m = new LinkedHashMap<>();
|
||||
m.put("state", s.state().name());
|
||||
m.put("detail", s.detail());
|
||||
return text(json(m));
|
||||
}
|
||||
|
||||
/** The one shared {@code NOT_CONFIGURED} refusal shape for {@code open}/{@code confirm}. */
|
||||
private static McpSchema.CallToolResult notConfigured() {
|
||||
return refusalJson(false, "NOT_CONFIGURED", "leadRollover: is not configured");
|
||||
@@ -2128,8 +2159,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 +2178,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"),
|
||||
@@ -2258,20 +2290,26 @@ public final class FleetMcp {
|
||||
return tool(FleetTool.HANDOVER.wireName(),
|
||||
"Replace your OWN lead session once its context is full: write a handover file, "
|
||||
+ "then use this to have fleetd clear your pane and bootstrap a fresh lead "
|
||||
+ "session against it. Three actions: 'open' (requests a token and the "
|
||||
+ "session against it. Four actions: 'open' (requests a token and the "
|
||||
+ "handoverPath you must write the handover file to before confirming), "
|
||||
+ "'confirm' (validates every gate and — only if every one passes — schedules "
|
||||
+ "the roll; it does NOT itself clear the pane, the roll runs once this call's "
|
||||
+ "own turn ends), and 'cancel' (drops a pending request without rolling). "
|
||||
+ "Primary-only. There is deliberately no terminal/session/leadTerminal "
|
||||
+ "parameter: the pane to roll is always resolved from YOUR OWN connection, "
|
||||
+ "never a value you pass, so you can only ever roll yourself — never another "
|
||||
+ "lead. Requires leadRollover: to be configured; when it is not, every action "
|
||||
+ "returns a clean refusal naming NOT_CONFIGURED instead of failing.",
|
||||
+ "own turn ends), 'cancel' (drops a pending request without rolling), and "
|
||||
+ "'status' (read-only: what happened to a token after 'confirm' — still "
|
||||
+ "running (approved but not finished yet), the roll completed, the calling "
|
||||
+ "turn never settled within turnSettleSeconds so no /clear was ever sent, or "
|
||||
+ "/clear itself never settled so bootstrapText was never sent; never "
|
||||
+ "schedules, cancels or retries anything). Primary-only. "
|
||||
+ "There is deliberately no terminal/session/leadTerminal parameter: the pane "
|
||||
+ "to roll is always resolved from YOUR OWN connection, never a value you "
|
||||
+ "pass, so you can only ever roll yourself — never another lead. Requires "
|
||||
+ "leadRollover: to be configured; when it is not, every action returns a "
|
||||
+ "clean refusal naming NOT_CONFIGURED instead of failing.",
|
||||
objectSchema(Map.of(
|
||||
"action", stringProp("\"open\", \"confirm\" or \"cancel\""),
|
||||
"action", stringProp("\"open\", \"confirm\", \"cancel\" or \"status\""),
|
||||
"reason", stringProp("Free-text audit note for \"open\" (optional, logged only)"),
|
||||
"token", stringProp("The token \"open\" returned — required for \"confirm\" and \"cancel\""),
|
||||
"token", stringProp("The token \"open\" returned — required for \"confirm\", "
|
||||
+ "\"cancel\" and \"status\""),
|
||||
"operatorConfirmed", Map.of("type", "boolean",
|
||||
"description", "For \"confirm\": your answer to \"has the human operator "
|
||||
+ "confirmed this wipe\" (default false; only consulted when "
|
||||
|
||||
@@ -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,84 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 2: {@link Fleetd#claudeCodeLauncher} is the factory that replaced {@code
|
||||
* main}'s inline {@code new ClaudeCodeLauncher(...)} call, whose 10th argument is the CB-596 {@code
|
||||
* memberCredentials} policy supplier ({@code () -> config.get().memberCredentials()}). Before this
|
||||
* ticket that argument was untestable wiring: replacing it with {@code () -> null} compiled with 0
|
||||
* errors and left every existing test green, since no existing test builds the exact object {@code
|
||||
* main} wires and then spawns it. {@code memberCredentials} being a lambda is never itself {@code
|
||||
* null}, so {@link dev.ltms.fleet.member.HerdrPeerLauncher#applyMemberCredentialPolicy} sees {@code
|
||||
* memberCredentials.get() == null} and silently shadows nothing — reopening the exact CB-592
|
||||
* exposure gap CB-596's policy closed (gitea issue #82).
|
||||
*
|
||||
* <p>This test drives the factory with a real {@link ConfigRef} carrying a {@code
|
||||
* memberCredentials:} block, spawns through the resulting launcher, and inspects what {@code
|
||||
* tab.create} actually carried — the same observable surface {@code ClaudeCodeLauncherTest}'s
|
||||
* {@code everyKnownNameNotAllowedIsShadowedWithTheSentinel} uses for the launcher's own credential
|
||||
* policy, applied here to prove {@code main}'s wiring reaches it.
|
||||
*/
|
||||
class FleetdClaudeCodeLauncherCredentialWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
ltms-local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: coder
|
||||
memberCredentials:
|
||||
policy: deny-by-default
|
||||
known:
|
||||
- GITEA_ACCESS_TOKEN
|
||||
""";
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, String> startEnv(FakeHerdr herdr) {
|
||||
return (Map<String, String>) ((Map<String, Object>) herdr.lastCall("tab.create").params()).get("env");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("main's memberCredentials wiring reaches ClaudeCodeLauncher: a known-but-not-allowed "
|
||||
+ "name is shadowed on spawn")
|
||||
void memberCredentialsWiringReachesClaudeCodeLauncher(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
ClaudeCodeLauncher launcher = Fleetd.claudeCodeLauncher(new AgentControl(herdr),
|
||||
new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
cfg.profiles(), cfg, config);
|
||||
launcher.spawn();
|
||||
|
||||
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
|
||||
assertNotNull(shadowed,
|
||||
"GITEA_ACCESS_TOKEN is 'known' but not 'allow'-ed in the loaded config — it must be "
|
||||
+ "explicitly shadowed on spawn; replacing the memberCredentials supplier with "
|
||||
+ "() -> null at the Fleetd.claudeCodeLauncher call site must fail this "
|
||||
+ "assertion, since a null policy shadows nothing");
|
||||
assertFalse(shadowed.isBlank(), "the overlay value must be non-blank");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,86 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: {@link Fleetd#exhaustedPatternLookup} is the factory that replaced {@code
|
||||
* main}'s inline lambda — resolve a herdr {@code target} to its session's profile, then to that
|
||||
* profile's live {@link LiveExhaustedPatterns#patternFor}. Same shape as {@link
|
||||
* Fleetd#worktreeBranchLookup} (which {@code FleetdWorktreeBranchLookupTest} pins the same way).
|
||||
*
|
||||
* <p>Before this ticket the lambda was built inline in {@code main} and untestable: replacing it
|
||||
* with {@code target -> null} — the exact shape of {@link ExhaustedPatternLookup#none()} — compiled
|
||||
* with 0 errors and left every existing test green. Per the ticket, this is the worst consequence
|
||||
* in the whole #589 sweep: a genuine usage-limit refusal would stop being classified as {@code
|
||||
* BACKEND_EXHAUSTED} and would be handed back to a waiting {@code fleet_send} as if it were real
|
||||
* completed work.
|
||||
*/
|
||||
class FleetdExhaustedPatternLookupWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: claude-opus-5
|
||||
exhaustedPattern: "usage limit"
|
||||
""";
|
||||
|
||||
private static MemberSession session(String terminal, String profile) {
|
||||
return new MemberSession("pane-" + terminal, terminal, profile, MemberRole.DEV,
|
||||
"/cwd", null, 0L, 0L, 0, MemberSession.State.READY, null, null);
|
||||
}
|
||||
|
||||
private static LiveExhaustedPatterns liveExhaustedPatterns(Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
return new LiveExhaustedPatterns(() -> ConfigRef.fixed(cfg).get().profiles());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a known target resolves through its session's profile to that profile's live pattern")
|
||||
void knownTargetResolvesThroughItsProfile(@TempDir Path dir) throws Exception {
|
||||
LiveExhaustedPatterns patterns = liveExhaustedPatterns(dir);
|
||||
ExhaustedPatternLookup lookup = Fleetd.exhaustedPatternLookup(
|
||||
() -> List.of(session("term1", "terra")), patterns);
|
||||
|
||||
Pattern resolved = lookup.patternFor("term1");
|
||||
|
||||
assertNotNull(resolved,
|
||||
"the lookup must resolve term1 -> profile 'terra' -> LiveExhaustedPatterns.patternFor("
|
||||
+ "'terra') — replacing the lambda body with 'target -> null' at the "
|
||||
+ "Fleetd.exhaustedPatternLookup call site must fail this assertion");
|
||||
assertTrue(resolved.matcher("the usage limit has been reached").find());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unknown target resolves to null, not a thrown exception")
|
||||
void unknownTargetResolvesToNull(@TempDir Path dir) throws Exception {
|
||||
LiveExhaustedPatterns patterns = liveExhaustedPatterns(dir);
|
||||
ExhaustedPatternLookup lookup = Fleetd.exhaustedPatternLookup(
|
||||
() -> List.of(session("term1", "terra")), patterns);
|
||||
|
||||
assertNull(lookup.patternFor("term_stranger"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: {@link Fleetd#forwardingExhaustionSink} is the factory that replaced
|
||||
* {@code main}'s inline {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)} (fleetd #175's
|
||||
* construction-order break: the adapters need a sink before {@code sessions} exists to build the
|
||||
* real one). Before this ticket that call site was untestable wiring: replacing the supplier
|
||||
* argument with a hardcoded {@code () -> ExhaustionSink.none()} compiled with 0 errors and left
|
||||
* every existing test green, because no test builds the object {@code main} actually wires and
|
||||
* then mutates the reference afterward — every existing {@code ExhaustionSink.forwardingTo} caller
|
||||
* in this codebase reads and writes the SAME reference within one test, so a hardcoded-none supplier
|
||||
* and a correctly-forwarding one are indistinguishable to them.
|
||||
*
|
||||
* <p>This test builds the reference, builds the forwarder from it, and only THEN repoints the
|
||||
* reference at a spy sink — the discriminating order fleetd #175's whole design depends on
|
||||
* ({@code exhaustionSinkRef} starts at {@code none()} and is repointed once {@code sessions}
|
||||
* exists). A forwarder that captured a fixed target at construction time (the inert form) can never
|
||||
* see that later repoint.
|
||||
*/
|
||||
class FleetdExhaustionSinkForwardingWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("the forwarder reads the reference live: repointing it AFTER construction is honoured")
|
||||
void forwarderReadsTheReferenceLiveNotAFixedTargetCapturedAtConstruction() {
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
ExhaustionSink forwarder = Fleetd.forwardingExhaustionSink(exhaustionSinkRef);
|
||||
|
||||
AtomicBoolean spyCalled = new AtomicBoolean(false);
|
||||
exhaustionSinkRef.set((target, reason, profile) -> spyCalled.set(true));
|
||||
|
||||
forwarder.onExhausted("term_x", "usage limit reached", "terra");
|
||||
|
||||
assertTrue(spyCalled.get(),
|
||||
"forwardingExhaustionSink must delegate to whatever exhaustionSinkRef currently "
|
||||
+ "holds — hardcoding the supplier to () -> ExhaustionSink.none() at the "
|
||||
+ "Fleetd.forwardingExhaustionSink call site must fail this assertion, "
|
||||
+ "since the spy set into the reference after construction would never run");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("before any repoint, the forwarder is inert — it starts at none(), not a crash")
|
||||
void beforeAnyRepointTheForwarderIsInert() {
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
ExhaustionSink forwarder = Fleetd.forwardingExhaustionSink(exhaustionSinkRef);
|
||||
|
||||
AtomicBoolean spyCalled = new AtomicBoolean(false);
|
||||
forwarder.onExhausted("term_x", "usage limit reached", "terra");
|
||||
|
||||
assertFalse(spyCalled.get(), "nothing was ever wired to be called here — this only pins "
|
||||
+ "that the factory does not throw before a real sink is published");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: {@link Fleetd#publishExhaustionSink} is the factory that replaced {@code
|
||||
* main}'s previously untested two-statement sequence — build the real {@link
|
||||
* Fleetd#exhaustionSink}, then {@code exhaustionSinkRef.set(exhaustionSink)}. {@link
|
||||
* Fleetd#exhaustionSink} itself is already pinned by {@code FleetdExhaustionSinkWarningTest} (its
|
||||
* log text) — what was NEVER pinned is the {@code .set(...)} call: {@code main} could replace it
|
||||
* with {@code exhaustionSinkRef.set(ExhaustionSink.none())} and compile with 0 errors, leaving
|
||||
* every existing test green, because {@link Fleetd#exhaustionSink}'s own tests build and call the
|
||||
* sink directly, never through the reference {@code main} publishes it into.
|
||||
*
|
||||
* <p>This test proves the PUBLISHED reference — not a freshly rebuilt sink — is the one that
|
||||
* actually quarantines a credential, by reading {@link BackendQuarantine#isQuarantined} after
|
||||
* calling {@code exhaustionSinkRef.get().onExhausted(...)}, the same object {@link
|
||||
* Fleetd#forwardingExhaustionSink} forwards to in production.
|
||||
*/
|
||||
class FleetdExhaustionSinkPublishWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: claude-opus-5
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""";
|
||||
|
||||
private static SessionManager emptyRosterSessions() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
FleetConfig.Profile dummy = new FleetConfig.Profile(
|
||||
"dummy", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(dummy.profile(), dummy), dummy.profile(), _ -> "tok");
|
||||
// Never acquires a session — publishExhaustionSink's built sink resolves target -> profile
|
||||
// via the profileHint fallback (fleetd #234), exactly like OpenCodeLauncher's real call
|
||||
// site does, so this never needs a populated roster.
|
||||
return new SessionManager(launcher);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the published reference actually quarantines — not a rebuilt-but-never-set sink")
|
||||
void publishedReferenceActuallyQuarantines(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
Map<String, String> reasonByCredential = new HashMap<>();
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
|
||||
Fleetd.publishExhaustionSink(exhaustionSinkRef, emptyRosterSessions(), config, quarantine,
|
||||
reasonByCredential, cfg);
|
||||
exhaustionSinkRef.get().onExhausted("term_x", "The usage limit has been reached", "terra");
|
||||
|
||||
assertTrue(quarantine.isQuarantined("terra"),
|
||||
"publishExhaustionSink must repoint exhaustionSinkRef at the REAL sink — "
|
||||
+ "replacing the .set(...) call with exhaustionSinkRef.set(ExhaustionSink.none()) "
|
||||
+ "at the Fleetd.publishExhaustionSink call site must fail this assertion, "
|
||||
+ "since none()'s onExhausted does nothing");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("before publishing, the reference is still inert — no quarantine, no crash")
|
||||
void beforePublishingTheReferenceIsStillInert(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
|
||||
exhaustionSinkRef.get().onExhausted("term_x", "The usage limit has been reached", "terra");
|
||||
|
||||
assertFalse(quarantine.isQuarantined("terra"),
|
||||
"nothing was published yet — this only pins the starting state the other test's "
|
||||
+ "assertion actually distinguishes from");
|
||||
}
|
||||
}
|
||||
@@ -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,72 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: {@link Fleetd#liveExhaustedPatterns} is the factory that replaced {@code
|
||||
* main}'s inline {@code new LiveExhaustedPatterns(() -> config.get().profiles())}. Before this
|
||||
* ticket, that supplier argument was untestable wiring: replacing it with a hardcoded {@code () ->
|
||||
* Map.of()} compiled with 0 errors and left every existing test green, because {@code
|
||||
* LiveExhaustedPatternsTest} builds its own instance directly with a hand-supplied map and never
|
||||
* goes through {@code main}'s call site.
|
||||
*
|
||||
* <p>Silently losing this wiring means every profile's {@code exhaustedPattern} stops being
|
||||
* recognised — {@link Fleetd#exhaustedPatternLookup} would never see a match, and a genuine
|
||||
* usage-limit refusal would be handed back to a waiting {@code fleet_send} as real completed work.
|
||||
*/
|
||||
class FleetdLiveExhaustedPatternsWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: claude-opus-5
|
||||
exhaustedPattern: "usage limit"
|
||||
gx:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
""";
|
||||
|
||||
private static ConfigRef loadConfig(Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
return ConfigRef.fixed(cfg);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile with a configured exhaustedPattern is armed, with a compiled matcher")
|
||||
void configuredProfileIsArmed(@TempDir Path dir) throws Exception {
|
||||
LiveExhaustedPatterns patterns = Fleetd.liveExhaustedPatterns(loadConfig(dir));
|
||||
|
||||
assertTrue(patterns.armed("terra"),
|
||||
"the config's live profiles() supplier must reach LiveExhaustedPatterns — hardcoding "
|
||||
+ "the supplier to () -> Map.of() at the Fleetd.liveExhaustedPatterns call "
|
||||
+ "site must fail this assertion");
|
||||
assertTrue(patterns.patternFor("terra").matcher("the usage limit has been reached").find());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile with no configured exhaustedPattern is not armed, but is still resolvable")
|
||||
void unconfiguredProfileIsNotArmed(@TempDir Path dir) throws Exception {
|
||||
LiveExhaustedPatterns patterns = Fleetd.liveExhaustedPatterns(loadConfig(dir));
|
||||
|
||||
assertFalse(patterns.armed("gx"), "'gx' has no exhaustedPattern configured");
|
||||
assertNull(patterns.patternFor("gx"));
|
||||
}
|
||||
}
|
||||
@@ -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,77 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.member.OpenCodeLauncher;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 2: {@link Fleetd#openCodeLauncher} is the factory that replaced {@code main}'s
|
||||
* inline {@code new OpenCodeLauncher(...)} call — the {@code opencode} counterpart to {@link
|
||||
* Fleetd#claudeCodeLauncher}, extracted for the identical reason. Its {@code memberCredentials}
|
||||
* argument is the same {@code () -> config.get().memberCredentials()} supplier; replacing it with
|
||||
* {@code () -> null} compiled with 0 errors and left every existing test green before this ticket,
|
||||
* reopening the same CB-592 exposure gap CB-596's policy closed.
|
||||
*
|
||||
* <p>Same observable surface as {@code OpenCodeLauncherTest}'s own {@code memberCredentials} tests:
|
||||
* spawn through the launcher {@code main} actually wires and inspect what {@code tab.create}
|
||||
* carried.
|
||||
*/
|
||||
class FleetdOpenCodeLauncherCredentialWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
gemini:
|
||||
kind: opencode
|
||||
model: google/gemini-2.5-pro
|
||||
memberCredentials:
|
||||
policy: deny-by-default
|
||||
known:
|
||||
- GITEA_ACCESS_TOKEN
|
||||
""";
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, String> startEnv(FakeHerdr herdr) {
|
||||
return (Map<String, String>) ((Map<String, Object>) herdr.lastCall("tab.create").params()).get("env");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("main's memberCredentials wiring reaches OpenCodeLauncher: a known-but-not-allowed "
|
||||
+ "name is shadowed on spawn")
|
||||
void memberCredentialsWiringReachesOpenCodeLauncher(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
OpenCodeLauncher launcher = Fleetd.openCodeLauncher(new AgentControl(herdr),
|
||||
new WorkspaceControl(herdr), cfg.profiles(), cfg, config, ExhaustionSink.none());
|
||||
launcher.spawn();
|
||||
|
||||
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
|
||||
assertNotNull(shadowed,
|
||||
"GITEA_ACCESS_TOKEN is 'known' but not 'allow'-ed in the loaded config — it must be "
|
||||
+ "explicitly shadowed on spawn; replacing the memberCredentials supplier with "
|
||||
+ "() -> null at the Fleetd.openCodeLauncher call site must fail this "
|
||||
+ "assertion, since a null policy shadows nothing");
|
||||
assertFalse(shadowed.isBlank(), "the overlay value must be non-blank");
|
||||
}
|
||||
}
|
||||
@@ -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");
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -16,6 +16,7 @@ import org.junit.jupiter.api.io.TempDir;
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Function;
|
||||
@@ -1004,6 +1005,324 @@ class LeadRolloverTest {
|
||||
}
|
||||
}
|
||||
|
||||
// ---- CB-... : LeadRollover#status makes the outcome of a confirmed roll readable -----------
|
||||
|
||||
@Test
|
||||
@DisplayName("[STATUS 1] after the calling turn never settles, status() reports "
|
||||
+ "TURN_NEVER_SETTLED for that token — this assertion could not even be written before "
|
||||
+ "status() existed")
|
||||
void statusReportsTurnNeverSettledAfterTheRollIsAbandoned() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("working"); // the calling lead's own pane — never goes idle in this test
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config =
|
||||
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 1 /*turnSettleSeconds*/, 20, "text");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "every synchronous gate should pass; the refusal happens "
|
||||
+ "only inside the deferred continuation, which this test's synchronous runner has "
|
||||
+ "already run to completion by the time confirm() returns");
|
||||
assertEquals(0, promptCallCount(herdr), "sanity: /clear was never sent");
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.TURN_NEVER_SETTLED, status.state());
|
||||
assertTrue(status.detail().contains("turnSettleSeconds"), "the detail must name the knob a "
|
||||
+ "reader needs to raise: " + status.detail());
|
||||
assertTrue(status.detail().toLowerCase().contains("no /clear"), "the detail must say plainly "
|
||||
+ "that no /clear was ever sent: " + status.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[STATUS 2] a roll that completes reports ROLLED for its token")
|
||||
void statusReportsRolledForACompletedRoll() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr(); // default idle — a full successful roll
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
|
||||
assertEquals(2, promptCallCount(herdr), "sanity: /clear then bootstrapText were both sent");
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.ROLLED, status.state());
|
||||
assertNotNull(status.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[STATUS 3] a roll where /clear never settles reports CLEAR_NEVER_SETTLED — "
|
||||
+ "distinct from both ROLLED and TURN_NEVER_SETTLED")
|
||||
void statusReportsClearNeverSettledDistinctFromTheOtherTwoStates() throws IOException {
|
||||
// Idle until /clear is sent, then permanently working — the SECOND wait never settles.
|
||||
FakeHerdr fake = new FakeHerdr();
|
||||
HerdrClient flipsAfterClear = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) throws HerdrException {
|
||||
JsonNode result = fake.call(method, params);
|
||||
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
|
||||
fake.agentStatus("working");
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
fake.close();
|
||||
}
|
||||
};
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config =
|
||||
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 20, 1 /*clearSettleSeconds*/, "boot text");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(flipsAfterClear, config, () -> clock.addAndGet(500));
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "every synchronous gate passes; the refusal is logged only, "
|
||||
+ "deep inside the deferred continuation");
|
||||
assertEquals(1, promptCallCount(fake), "sanity: only /clear was sent, never bootstrapText");
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.CLEAR_NEVER_SETTLED, status.state());
|
||||
assertNotEquals(LeadRollover.RollState.ROLLED, status.state());
|
||||
assertNotEquals(LeadRollover.RollState.TURN_NEVER_SETTLED, status.state());
|
||||
assertTrue(status.detail().contains("clearSettleSeconds"), status.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[STATUS 4] a token that was never issued, or was cancelled, gives a clean "
|
||||
+ "UNKNOWN answer rather than an exception or a false ROLLED")
|
||||
void statusOnUnissuedOrCancelledTokenIsCleanNotAnExceptionOrFalseSuccess() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
|
||||
|
||||
// never issued at all
|
||||
LeadRollover.RollStatus neverIssued = assertDoesNotThrow(() -> rollover.status("no-such-token"));
|
||||
assertEquals(LeadRollover.RollState.UNKNOWN, neverIssued.state());
|
||||
assertNotEquals(LeadRollover.RollState.ROLLED, neverIssued.state());
|
||||
|
||||
// null / blank must not throw either — pending is a ConcurrentHashMap, which throws on a
|
||||
// null-key lookup unless status() guards it first
|
||||
assertEquals(LeadRollover.RollState.UNKNOWN, assertDoesNotThrow(() -> rollover.status(null)).state());
|
||||
assertEquals(LeadRollover.RollState.UNKNOWN, assertDoesNotThrow(() -> rollover.status(" ")).state());
|
||||
|
||||
// opened, then cancelled
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
assertTrue(rollover.cancel(pending.token()));
|
||||
LeadRollover.RollStatus cancelled = assertDoesNotThrow(() -> rollover.status(pending.token()));
|
||||
assertEquals(LeadRollover.RollState.UNKNOWN, cancelled.state());
|
||||
assertNotEquals(LeadRollover.RollState.ROLLED, cancelled.state());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[STATUS 5] the bounded outcome history never grows past its cap")
|
||||
void statusHistoryDoesNotGrowPastItsCap() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr(); // default idle — every roll completes
|
||||
Path handover = writeHandover("handover contents"); // one file, reused by every roll below —
|
||||
// checkHandover only compares its mtime against each open()'s OWN requestedAtMillis (the
|
||||
// fake clock, in the low thousands), and the file's real (wall-clock) mtime is always far
|
||||
// larger than that, so freshness passes on every iteration without rewriting the file.
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), () -> clock.addAndGet(1));
|
||||
|
||||
int rolls = LeadRollover.OUTCOME_HISTORY_CAP + 50;
|
||||
String[] tokens = new String[rolls];
|
||||
for (int i = 0; i < rolls; i++) {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full #" + i);
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "roll #" + i + " should have been approved: " + decision.detail());
|
||||
tokens[i] = pending.token();
|
||||
}
|
||||
|
||||
assertEquals(LeadRollover.RollState.ROLLED, rollover.status(tokens[rolls - 1]).state(),
|
||||
"the most recently finished roll's outcome must still be in the bounded history");
|
||||
|
||||
// There is no direct size accessor for the bounded history, so boundedness is asserted
|
||||
// indirectly and behaviourally: the OLDEST finished roll's outcome must have been evicted
|
||||
// (reads back as a clean UNKNOWN, exactly like a token that was never issued) once more than
|
||||
// OUTCOME_HISTORY_CAP rolls have gone through this instance. If the cap were not enforced,
|
||||
// tokens[0] would still read back ROLLED here, and this assertion would fail.
|
||||
assertEquals(LeadRollover.RollState.UNKNOWN, rollover.status(tokens[0]).state(),
|
||||
"the oldest finished roll's outcome must have been evicted once the cap was "
|
||||
+ "exceeded — otherwise the bounded history is not actually bounded");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[STATUS 6] a status() call against a pending roll is read-only — no herdr call is "
|
||||
+ "made and the pending request is left untouched")
|
||||
void statusCallAgainstAPendingRollIsReadOnlyAndTouchesNothing() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.PENDING, status.state());
|
||||
assertEquals(0, herdr.calls.size(), "status() must never make any herdr call at all — not "
|
||||
+ "just no agent.prompt — since it must never schedule, cancel, or retry anything");
|
||||
|
||||
// the pending request must be left exactly as it was: the SAME token can still be confirmed
|
||||
// afterwards, as if status() had never been called.
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "status() must not have consumed or otherwise disturbed the "
|
||||
+ "pending request: " + decision.reason() + " / " + decision.detail());
|
||||
}
|
||||
|
||||
// ---- PR #600 review round 2: IN_PROGRESS — the gap between confirm() handing off and the ----
|
||||
// ---- continuation finishing must never read back as UNKNOWN ("nothing was ever requested") --
|
||||
|
||||
/**
|
||||
* A {@code continuationRunner} that CAPTURES the roll instead of running it, so a test can
|
||||
* observe {@link LeadRollover#status} in the window between {@link LeadRollover#confirm}
|
||||
* handing off and the roll actually finishing — the window the synchronous {@code
|
||||
* Runnable::run} runner used everywhere else in this class collapses to nothing. Call {@link
|
||||
* #runNext()} to finish exactly one held roll, once the test is done observing the in-flight
|
||||
* state.
|
||||
*/
|
||||
private static final class HoldingRunner implements java.util.function.Consumer<Runnable> {
|
||||
private final List<Runnable> held = new ArrayList<>();
|
||||
|
||||
@Override
|
||||
public void accept(Runnable runnable) {
|
||||
held.add(runnable);
|
||||
}
|
||||
|
||||
int heldCount() {
|
||||
return held.size();
|
||||
}
|
||||
|
||||
/** Runs (and removes) the oldest held roll — FIFO, matching confirm() call order. */
|
||||
void runNext() {
|
||||
held.remove(0).run();
|
||||
}
|
||||
}
|
||||
|
||||
private static LeadRollover newRolloverWithHoldingRunner(HerdrClient herdr,
|
||||
FleetConfig.LeadRollover config, LongSupplier nowMillis, HoldingRunner runner) {
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
return new LeadRollover(agents, () -> config, _ -> null, nowMillis, () -> { }, runner);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[IN-PROGRESS 1] between an approved confirm() and the continuation finishing, "
|
||||
+ "status() reports IN_PROGRESS — not UNKNOWN, and not PENDING")
|
||||
void statusReportsInProgressBetweenConfirmAndTheContinuationFinishing() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr(); // default idle — the held roll WOULD complete once run
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
HoldingRunner runner = new HoldingRunner();
|
||||
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
|
||||
fixedClock(clock), runner);
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
|
||||
assertEquals(1, runner.heldCount(), "sanity: the roll must have been handed to the "
|
||||
+ "continuation runner and held there, not run yet");
|
||||
assertEquals(0, promptCallCount(herdr), "sanity: the held continuation has not run, so "
|
||||
+ "nothing has been sent to the pane yet");
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.IN_PROGRESS, status.state(),
|
||||
"a lead calling status() right after confirm() returned approved, while the roll "
|
||||
+ "is still running, must be told IN_PROGRESS — not UNKNOWN (\"nothing was "
|
||||
+ "ever requested\", which would wrongly invite it to call open() again "
|
||||
+ "mid-roll) and not PENDING (\"not yet approved\", which is simply false "
|
||||
+ "here): got " + status.state() + " / " + status.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[IN-PROGRESS 2] once the continuation finishes, the same token reports its "
|
||||
+ "terminal state — IN_PROGRESS is not sticky")
|
||||
void inProgressStateIsNotStickyOnceTheContinuationFinishes() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr(); // default idle — the held roll completes once run
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
HoldingRunner runner = new HoldingRunner();
|
||||
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
|
||||
fixedClock(clock), runner);
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
|
||||
assertEquals(LeadRollover.RollState.IN_PROGRESS, rollover.status(pending.token()).state(),
|
||||
"sanity: must be IN_PROGRESS before the held continuation is run");
|
||||
|
||||
runner.runNext(); // finish the held roll now
|
||||
|
||||
assertEquals(2, promptCallCount(herdr), "sanity: the roll actually ran to completion "
|
||||
+ "once released — /clear then bootstrapText");
|
||||
LeadRollover.RollStatus after = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.ROLLED, after.state(),
|
||||
"the SAME token must now report its terminal state — IN_PROGRESS must not still be "
|
||||
+ "reported once the roll has actually finished");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[IN-PROGRESS 4] IN_PROGRESS is distinct from every other RollState")
|
||||
void inProgressStateIsDistinctFromAllOtherStates() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
HoldingRunner runner = new HoldingRunner();
|
||||
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
|
||||
fixedClock(clock), runner);
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
|
||||
|
||||
LeadRollover.RollState inProgress = rollover.status(pending.token()).state();
|
||||
assertEquals(LeadRollover.RollState.IN_PROGRESS, inProgress);
|
||||
for (LeadRollover.RollState other : LeadRollover.RollState.values()) {
|
||||
if (other == LeadRollover.RollState.IN_PROGRESS) {
|
||||
continue;
|
||||
}
|
||||
assertNotEquals(other, inProgress, "IN_PROGRESS must be distinct from " + other);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[IN-PROGRESS 5] eviction counts IN_PROGRESS entries toward the cap exactly like "
|
||||
+ "finished ones — a burst of confirmed-but-not-yet-finished rolls still ages out")
|
||||
void evictionCountsInProgressEntriesTowardTheCap() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents"); // reused by every roll — see
|
||||
// statusHistoryDoesNotGrowPastItsCap for why one shared file is enough for freshness.
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
HoldingRunner runner = new HoldingRunner(); // nothing run below — every roll stays IN_PROGRESS
|
||||
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
|
||||
() -> clock.addAndGet(1), runner);
|
||||
|
||||
int rolls = LeadRollover.OUTCOME_HISTORY_CAP + 50;
|
||||
String[] tokens = new String[rolls];
|
||||
for (int i = 0; i < rolls; i++) {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full #" + i);
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "roll #" + i + " should have been approved: " + decision.detail());
|
||||
tokens[i] = pending.token();
|
||||
}
|
||||
assertEquals(rolls, runner.heldCount(), "sanity: none of these rolls have been run — every "
|
||||
+ "one of them is sitting in outcomes as IN_PROGRESS, not a separate uncapped map");
|
||||
|
||||
assertEquals(LeadRollover.RollState.IN_PROGRESS, rollover.status(tokens[rolls - 1]).state(),
|
||||
"the most recently confirmed (still in-flight) roll must still be in the bounded "
|
||||
+ "history");
|
||||
assertEquals(LeadRollover.RollState.UNKNOWN, rollover.status(tokens[0]).state(),
|
||||
"the oldest confirmed roll's IN_PROGRESS entry must have been evicted once the cap "
|
||||
+ "was exceeded, exactly like a finished entry would be — proving IN_PROGRESS "
|
||||
+ "entries share the SAME bounded map and count against the SAME cap, rather "
|
||||
+ "than living in a second, uncapped in-flight map");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494 follow-up] the turn-settle timeout warn line prints the MEASURED "
|
||||
+ "elapsed time next to the configured budget, never the configured value alone")
|
||||
|
||||
@@ -258,5 +258,51 @@ class FleetMcpHandoverTest {
|
||||
assertTrue(textOf(r).contains("\"cancelled\":true"), textOf(r));
|
||||
}
|
||||
|
||||
// --- new: action "status" — makes the outcome of a confirm() readable through the tool ----
|
||||
|
||||
@Test
|
||||
@DisplayName("status on a null LeadRollover is a clean NOT_CONFIGURED refusal, never a throw")
|
||||
void statusWithNullLeadRolloverRefusesCleanly() {
|
||||
McpSchema.CallToolResult r = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
|
||||
Map.of("action", "status", "token", "whatever")));
|
||||
assertFalse(r.isError());
|
||||
assertTrue(textOf(r).contains("NOT_CONFIGURED"), textOf(r));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("status on a token that was never opened reports UNKNOWN")
|
||||
void statusOnUnknownTokenReportsUnknown() {
|
||||
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
|
||||
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
|
||||
Map.of("action", "status", "token", "does-not-exist"));
|
||||
assertFalse(r.isError());
|
||||
assertTrue(textOf(r).contains("\"state\":\"UNKNOWN\""), textOf(r));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("status on a token that is still pending (opened, not confirmed) reports PENDING")
|
||||
void statusOnPendingTokenReportsPending() {
|
||||
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
|
||||
String token = extractToken(textOf(FleetMcp.handover(rollover, LEAD, Map.of("action", "open"))));
|
||||
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
|
||||
Map.of("action", "status", "token", token));
|
||||
assertFalse(r.isError());
|
||||
assertTrue(textOf(r).contains("\"state\":\"PENDING\""), textOf(r));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("status is registered on the tool's schema and the schema still names no caller-identity parameter")
|
||||
void statusActionIsAdvertisedOnTheSchema() {
|
||||
FleetMcp m = mcp(null);
|
||||
McpSchema.Tool tool = m.registeredTools().stream()
|
||||
.filter(t -> "fleet_handover".equals(t.name()))
|
||||
.findFirst()
|
||||
.orElseThrow(() -> new AssertionError("fleet_handover was not registered"));
|
||||
assertTrue(tool.description().contains("'status'"),
|
||||
"the tool's own description must advertise the 'status' action: " + tool.description());
|
||||
}
|
||||
|
||||
// --- acceptance 7 (wiring) is covered by FleetdLeadRolloverWiringTest, unchanged -----------
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
|
||||
+131
-23
@@ -50,6 +50,14 @@
|
||||
# absence is the real signal, because a drain that dies on its first session prints nothing
|
||||
# else either. Warns loudly; never fails the redeploy, because by the time this is detectable
|
||||
# the new daemon is already up and healthy.
|
||||
# 10. fleetd #603 — the same shape as trap 3 above, through a different door: the step that waited
|
||||
# for the NEW process to appear gave it its own short, fixed 10s budget, then hard-`die`d,
|
||||
# while the health check right after it waits a full $HEALTH_WAIT (60s) for the same daemon to
|
||||
# answer. Under launchd, `launchctl load` returns as soon as launchd accepts the job, before the
|
||||
# java process exists, and on a slow host that took longer than 10s — so the script died with
|
||||
# "no process appeared" on a deploy that had fully succeeded. The pid poll now shares
|
||||
# $HEALTH_WAIT instead of a separate, shorter budget, and a miss there falls through to the
|
||||
# health check (the truer signal: is it actually answering?) instead of killing the run.
|
||||
#
|
||||
# Usage:
|
||||
# scripts/redeploy-fleetd.sh # build, confirm, restart, verify
|
||||
@@ -79,7 +87,9 @@ OUT="$MODULE/fleetd.out"
|
||||
PATTERN='target/fleetd.jar'
|
||||
HEALTH='http://127.0.0.1:8765/healthz'
|
||||
STOP_WAIT=30 # seconds to wait for a clean exit before reporting failure
|
||||
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start
|
||||
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start — fleetd #603: also the pid-
|
||||
# poll budget below (wait_for_new_pid/await_daemon_started), so the two checks
|
||||
# share one named budget instead of the pid poll holding its own shorter one
|
||||
|
||||
# CB-594: the launchd agent this script must not fight with (see trap 6 above).
|
||||
LAUNCHD_LABEL='dev.ltms.fleetd'
|
||||
@@ -169,7 +179,55 @@ hash256() {
|
||||
# on PATH). "absent" must never be the answer for a file that exists — that conflation, on Linux,
|
||||
# was the whole defect this ticket fixes.
|
||||
jar_id() { local f="${1:-$JAR}"; [ -f "$f" ] && hash256 "$f" || echo "absent"; }
|
||||
running_pid() { pgrep -f "$PATTERN" || true; }
|
||||
|
||||
# fleetd #593 — `pgrep -f "$PATTERN"` matches ANY process whose full command line CONTAINS the
|
||||
# pattern text, and that is not the same thing as "is the daemon". A shell that merely embeds the
|
||||
# pattern as literal text — a human typing this exact investigation by hand, an ssh-shaped
|
||||
# `sh -c '...; ...'`, a pipeline, or any other non-exec'ing shell that never replaced itself with
|
||||
# the pattern-holding command — still shows up in that match, and it is the INSTRUMENT, not the
|
||||
# daemon. Measured live on this Mac: `sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 30' &`
|
||||
# leaves a real `sh` process alive (it forks for the `sleep`, it does not exec into it) whose own
|
||||
# `ps -o args` is `sh -c echo "target/fleetd.jar" >/dev/null; sleep 30` — `pgrep -f "$PATTERN"`
|
||||
# matches that line right alongside the real `java -jar target/fleetd.jar` process. `pgrep -c`
|
||||
# (an in-one-call count) does not exist on BSD/macOS at all, so this cannot be fixed by switching
|
||||
# pgrep flags — it has to filter what pgrep already found, after the fact, in a way that still
|
||||
# runs on BSD.
|
||||
#
|
||||
# fleetd #593 CORRECTION 1 — the first cut of this filter kept everything whose `comm` was NOT a
|
||||
# shell name (a denylist: sh/bash/zsh/dash/ksh). Two holes in that, both the same false-positive
|
||||
# shape the ticket exists to remove in the first place:
|
||||
# 1. a pid `pgrep` just listed can exit before the `ps -o comm=` lookup runs; on a gone pid `ps`
|
||||
# prints nothing, `comm` ends up empty, and an empty string matches none of the denied shell
|
||||
# names — so a pid that no longer exists was still counted.
|
||||
# 2. the denylist only knows the shells someone thought to name. `ssh`, `perl`, `python3`,
|
||||
# `ruby`, `tail` — anything else that carries the pattern in its own argv — was still
|
||||
# counted right along with the real daemon, and the ticket names `ssh` as a live route.
|
||||
# Both close with the same change: allowlist `comm = java` instead of denying shells. Measured on
|
||||
# the live daemon: `pid=30224 comm=java`. An empty comm (hole 1) is not `java` either, so it is
|
||||
# excluded for free — no separate "is this pid still alive" check needed.
|
||||
#
|
||||
# The objection, because it is real: an allowlist can UNDER-count. If fleetd ever stops being
|
||||
# launched as `java -jar ...` — a native image, a renamed launcher — `running_pid()` silently
|
||||
# returns nothing and `assert_single_daemon` stops noticing a second daemon at all. For a guard,
|
||||
# that false-negative direction is the worse one to be wrong in. This is not a new assumption,
|
||||
# though: `PATTERN='target/fleetd.jar'` two lines up already assumes the daemon is a jar, which
|
||||
# is only ever run by `java`. If that launch method changes, `PATTERN` stops matching anything
|
||||
# before this allowlist would ever get the chance to be wrong — the allowlist rides on the same
|
||||
# assumption that is already load-bearing, it does not add a new one. Whoever changes the launch
|
||||
# method needs to update both `PATTERN` and this allowlist together.
|
||||
running_pid() {
|
||||
local pid comm out=''
|
||||
for pid in $(pgrep -f "$PATTERN" 2>/dev/null || true); do
|
||||
comm="$(ps -o comm= -p "$pid" 2>/dev/null || true)"
|
||||
comm="${comm##*/}"
|
||||
comm="${comm#-}"
|
||||
# Allowlist, not a denylist of wrappers — see the CORRECTION 1 comment above. Anything that
|
||||
# is not literally `java` is excluded, including an empty comm from a pid that already exited.
|
||||
[ "$comm" = java ] || continue
|
||||
out="$out$pid"$'\n'
|
||||
done
|
||||
printf '%s' "$out"
|
||||
}
|
||||
|
||||
# fleetd #493 — three small, independently testable pieces of "never build into the path a
|
||||
# running process holds":
|
||||
@@ -507,8 +565,11 @@ assert_single_daemon() {
|
||||
die "more than one fleetd process is running after this restart (pids: $(printf '%s' "$pids" | tr '\n' ' ')).
|
||||
This is the exact failure a racing supervisor produces: the OLD jar was revived by its
|
||||
supervisor while this script started a NEW copy. Two daemons on one herdr session kill
|
||||
each other's members. Investigate with 'pgrep -f \"$PATTERN\"' and stop the wrong one by
|
||||
hand — do not assume either pid is the one you want."
|
||||
each other's members. Investigate with 'ps -eo pid,comm,args | grep -F \"$PATTERN\"' and
|
||||
check the COMM column of each hit yourself before acting — a bare 'pgrep -f \"$PATTERN\"'
|
||||
(fleetd #593) can match the very shell you type it into, not just the daemon, so it is not
|
||||
safe remediation advice on its own. Stop the wrong one by hand — do not assume either pid
|
||||
is the one you want."
|
||||
fi
|
||||
}
|
||||
|
||||
@@ -1016,6 +1077,67 @@ $(tail -30 "$out_file" 2>/dev/null)"
|
||||
fi
|
||||
}
|
||||
|
||||
# fleetd #603 — the dual of wait_for_daemon_exit above: waits for a pid to APPEAR instead of
|
||||
# disappear. Used to be a bare `for _ in $(seq 10)` sitting directly in the main flow, with its own
|
||||
# short, fixed budget that had nothing to do with $HEALTH_WAIT (60s) — the budget the health check
|
||||
# right after it gets for the very same daemon. Under launchd, `launchctl load` returns as soon as
|
||||
# launchd accepts the job, before the java process exists, and on a slow host that took longer than
|
||||
# 10s — so the script died with "no process appeared" on a deploy that had fully succeeded (the
|
||||
# operator confirmed pid, healthz, and a fresh log line all present, by hand, right afterwards).
|
||||
# Sharing $HEALTH_WAIT here removes the extra, shorter magic number without inventing a new one.
|
||||
wait_for_new_pid() {
|
||||
local timeout="$1" _i
|
||||
for _i in $(seq "$timeout"); do
|
||||
[ -n "$(running_pid)" ] && return 0
|
||||
sleep 1
|
||||
done
|
||||
[ -n "$(running_pid)" ]
|
||||
}
|
||||
|
||||
# fleetd #603 — the pid-appeared check and the healthz check, folded into one decision. Same shape
|
||||
# as swap_if_built/refuse_drain_gate/report_shutdown_drain above (#521/#528/#512): the main flow
|
||||
# calls this ONE function unconditionally, so there is no bare guard left for a future edit to
|
||||
# invert independently of it. That matters more here than for most of those: a plain grep of this
|
||||
# script's source cannot tell "a pid miss falls through to the health check" from "a pid miss still
|
||||
# dies" apart, because both read as the same two lines of text with only the runtime branch
|
||||
# changed — a source-text test could pass on either behavior. A real behavioural test on this
|
||||
# function is the only thing that can actually tell them apart, which is why one exists below.
|
||||
#
|
||||
# wait_for_new_pid's miss is no longer fatal by itself: it falls through to the health check, which
|
||||
# is direct proof the new daemon is up (/healthz answers 200) rather than a proxy for it (a process
|
||||
# merely existing under a name running_pid() recognises). A genuine failure still dies here: it
|
||||
# misses the pid poll AND the health poll, and report_health's own die() still prints the tail of
|
||||
# $out_file, exactly as before this fix.
|
||||
#
|
||||
# Sets NEW_PID (global — the caller's "pid N, jar ..." result line reads it afterwards) and
|
||||
# HEALTH_BODY/HEALTH_CODE (globals, the same reason report_health already needs them handed back).
|
||||
# Only ever called as a bare statement in the main flow below, never from inside a `$( )`: a die()
|
||||
# reached from inside a command substitution only kills that subshell, not the whole script, which
|
||||
# would silently turn a genuine failure back into a false "succeeded" exit (see poll_health_body's
|
||||
# own `|| true` idiom for the same hazard from the other direction).
|
||||
await_daemon_started() {
|
||||
local health_wait="$1" old_pid="$2" health_url="$3" out_file="$4"
|
||||
NEW_PID=""
|
||||
if wait_for_new_pid "$health_wait"; then
|
||||
NEW_PID="$(running_pid)"
|
||||
[ "$NEW_PID" != "${old_pid:-}" ] || die "pid unchanged ($NEW_PID) — the old daemon never died"
|
||||
ok "started, pid $NEW_PID"
|
||||
else
|
||||
warn "no process matched $PATTERN within ${health_wait}s of starting — falling through to the health check, which is the more truthful signal"
|
||||
fi
|
||||
|
||||
HEALTH_BODY="$(poll_health_body "$health_url" "$health_wait")" || true
|
||||
HEALTH_CODE="000"
|
||||
[ -n "$HEALTH_BODY" ] || HEALTH_CODE="$(curl -s -o /dev/null -w '%{http_code}' --max-time 2 "$health_url" 2>/dev/null || echo 000)"
|
||||
report_health "$HEALTH_BODY" "$HEALTH_CODE" "$out_file" "$health_wait"
|
||||
|
||||
# By now /healthz has answered (report_health above would have died otherwise), so the daemon is
|
||||
# confirmed up even if the pid poll never matched it — see running_pid()'s own comment on
|
||||
# under-counting if the launch method ever stops being a plain `java -jar`. Fill NEW_PID in for the
|
||||
# result line rather than leave it blank on an otherwise fully successful redeploy.
|
||||
[ -n "$NEW_PID" ] || NEW_PID="$(running_pid)"
|
||||
}
|
||||
|
||||
# Ticket item 8 — `HAD_OLD_PID=0; [ -n "$OLD_PID" ] && HAD_OLD_PID=1`. Measured safe under `set -e`
|
||||
# at both bash 3.2.57 and 5.x (see the header comment trap 9 discussion in the ticket) — not a `set
|
||||
# -e` hazard, but still an untested computation feeding report_shutdown_drain's own four-way
|
||||
@@ -1234,28 +1356,14 @@ say "start"
|
||||
# there now, so this line is the only thing left in the main flow to get wrong.
|
||||
dispatch_start "$SUPERVISOR_KIND"
|
||||
|
||||
for _ in $(seq 10); do
|
||||
NEW_PID="$(running_pid)"
|
||||
[ -n "$NEW_PID" ] && break
|
||||
sleep 1
|
||||
done
|
||||
[ -n "${NEW_PID:-}" ] || die "no process appeared. Last lines of $OUT:
|
||||
$(tail -20 "$OUT" 2>/dev/null)"
|
||||
[ "$NEW_PID" != "${OLD_PID:-}" ] || die "pid unchanged ($NEW_PID) — the old daemon never died"
|
||||
ok "started, pid $NEW_PID"
|
||||
|
||||
# ------------------------------------------------------------------ verify
|
||||
|
||||
say "verify"
|
||||
|
||||
# fleetd #555: poll_health_body/health_is_up/report_health above. HEALTH_CODE is only ever
|
||||
# consulted by report_health when the body came back empty; `|| true` on both assignments is the
|
||||
# same "an absent/failing command substitution must not kill the script under set -e" idiom the
|
||||
# swap/drain helpers already rely on (see poll_health_body's own comment).
|
||||
HEALTH_BODY="$(poll_health_body "$HEALTH" "$HEALTH_WAIT")" || true
|
||||
HEALTH_CODE="000"
|
||||
[ -n "$HEALTH_BODY" ] || HEALTH_CODE="$(curl -s -o /dev/null -w '%{http_code}' --max-time 2 "$HEALTH" 2>/dev/null || echo 000)"
|
||||
report_health "$HEALTH_BODY" "$HEALTH_CODE" "$OUT" "$HEALTH_WAIT"
|
||||
# fleetd #603: await_daemon_started above folds the pid-appeared check and the healthz check into
|
||||
# one decision — see its own comment for why a bare `if` here, split back into the old two pieces,
|
||||
# would put back an untestable branch this ticket exists to close.
|
||||
await_daemon_started "$HEALTH_WAIT" "$OLD_PID" "$HEALTH" "$OUT"
|
||||
|
||||
# A fresh listening line, strictly after the restart mark. An old daemon that never died would
|
||||
# otherwise let an old line pass for a new one.
|
||||
@@ -1310,7 +1418,7 @@ report_shutdown_drain "$FRESH_LOG" "$HAD_OLD_PID"
|
||||
assert_single_daemon "$(running_pid)"
|
||||
|
||||
say "result"
|
||||
ok "pid $NEW_PID, jar $(jar_id)"
|
||||
ok "pid ${NEW_PID:-unknown}, jar $(jar_id)"
|
||||
if [ "$REDEPLOY_AMQP_CHECK_SKIPPED" -eq 1 ]; then
|
||||
# fleetd #552: the fourth reader of the fresh-log region. Without this branch REDEPLOY_ERROR_COUNT
|
||||
# stays at its untouched 0 (classify_amqp_connection_errors never got a file to read) and
|
||||
|
||||
@@ -433,6 +433,130 @@ test_assert_single_daemon_rejects_two_pids() {
|
||||
printf '%s' "$output" | grep -qF '4343' || fail "refusal message does not list the pids it found"
|
||||
}
|
||||
|
||||
# fleetd #593 instance 2 — `running_pid()` used to be a bare `pgrep -f "$PATTERN"`, which matches
|
||||
# ANY process whose full command line contains the pattern TEXT, including a shell that merely
|
||||
# embeds it as literal text rather than being the daemon. Measured live on this Mac: `pgrep -c`
|
||||
# (a one-call count) does not exist on BSD at all, and `bash -c "<single command>"` execs in place
|
||||
# so no parent shell survives to hold the pattern — which is exactly why the defect did not
|
||||
# reproduce from a plain script and needs a wrapper shaped like this instead. A `sh -c '...; ...'`
|
||||
# with MORE THAN ONE statement does not get that exec-in-place treatment: the shell forks a child
|
||||
# for the second statement and stays alive itself, holding the whole `-c` string — pattern text
|
||||
# included — in its own `ps -o args`, for as long as it runs. That is the same shape an
|
||||
# `ssh host "…; …"` wrapper or a hand-typed pipeline leaves behind. Before the fix this test would
|
||||
# have found the wrapper's pid in running_pid()'s output; it must not.
|
||||
test_running_pid_excludes_self_matching_wrapper_shell() {
|
||||
local before after wrapper_pid
|
||||
before="$(running_pid)"
|
||||
sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' &
|
||||
wrapper_pid=$!
|
||||
sleep 0.3
|
||||
after="$(running_pid)"
|
||||
kill "$wrapper_pid" 2>/dev/null || true
|
||||
wait "$wrapper_pid" 2>/dev/null || true
|
||||
[ "$after" = "$before" ] \
|
||||
|| fail "running_pid() counted a self-matching wrapper shell (pid $wrapper_pid, holding the pattern as literal text in its own argv, not the daemon): before=[$before] after=[$after]"
|
||||
}
|
||||
|
||||
# fleetd #593 CORRECTION 1 — the round-1 version of this test gave its standin an argv[0]
|
||||
# containing the pattern text (via `exec -a`) and left `comm` as whatever that override produced,
|
||||
# which was never `java`. That was fine for a denylist-of-shells filter, but the allowlist below
|
||||
# now requires `comm = java` specifically, so the standin here must actually carry that comm, not
|
||||
# just avoid being a shell. `exec -a java` overrides argv[0] to `java` while the process itself
|
||||
# stays a genuine, harmless `sh`; combining it with the same non-exec'ing multi-statement shape
|
||||
# the wrapper-shell test above uses keeps the pattern text in the process's own `ps -o args` for
|
||||
# as long as it runs. Measured live on this Mac (BSD/macOS: `ps -o comm=` here reflects argv[0]):
|
||||
# `comm=java`, `args` contains the pattern, `pgrep -f "$PATTERN"` finds it. Copying a real system
|
||||
# binary into a scratch path and executing it from there was tried first, for a more literal
|
||||
# stand-in daemon, and the OS killed it outright (SIGKILL, exit 137 — almost certainly a
|
||||
# code-signing check on a relocated binary); `exec -a` needs no binary of its own and nothing
|
||||
# under a scratch directory, and it is the technique CORRECTION 1 names as the right one.
|
||||
#
|
||||
# This is the one live-process test in this file whose result could differ on Linux: Linux sets
|
||||
# `comm` from the actually-executed binary's own path, not from `exec -a`'s argv[0] override (BSD
|
||||
# ties `comm` to argv[0], which is what makes this technique work here) — so on Linux this
|
||||
# specific fixture might report `comm=sh`, not `comm=java`, even though the REAL daemon (a literal
|
||||
# `java -jar target/fleetd.jar` process, never fabricated) is unaffected either way. I could not
|
||||
# verify this fixture's behavior on Linux, so test_running_pid_counts_a_pid_whose_comm_is_java
|
||||
# below backstops the same claim (the allowlist admits a pid whose comm is `java`) with a stubbed
|
||||
# `ps`, which is identical bash on every platform and carries no such platform question.
|
||||
test_running_pid_finds_a_real_java_named_second_process() {
|
||||
local before after standin_pid
|
||||
before="$(running_pid)"
|
||||
( exec -a java sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' ) &
|
||||
standin_pid=$!
|
||||
sleep 0.3
|
||||
after="$(running_pid)"
|
||||
kill "$standin_pid" 2>/dev/null || true
|
||||
wait "$standin_pid" 2>/dev/null || true
|
||||
printf '%s\n' "$after" | grep -qxF "$standin_pid" \
|
||||
|| fail "running_pid() did not find a real second process (pid $standin_pid, comm forced to 'java' via exec -a) whose own argv holds the pattern: before=[$before] after=[$after]"
|
||||
}
|
||||
|
||||
# fleetd #593 CORRECTION 1, hole 2 — the round-1 filter denied known shell names (sh/bash/zsh/
|
||||
# dash/ksh) and counted everything else. `ssh`, `perl`, `python3`, `ruby`, `tail` — anything not on
|
||||
# that list, carrying the pattern in its own argv — was still counted right alongside the real
|
||||
# daemon, and the ticket names `ssh` as a live route. Stubbing `pgrep`/`ps` (rather than spawning a
|
||||
# real perl/ssh process) pins the exact discriminator this correction is about — comm, not the
|
||||
# caller's shape — deterministically on every platform, with no dependency on perl/python3/ruby
|
||||
# being installed in whatever environment runs this suite, and no dependency on how a given OS
|
||||
# derives `comm` for a fabricated process (see the comment above
|
||||
# test_running_pid_finds_a_real_java_named_second_process for why that matters here).
|
||||
test_running_pid_drops_a_pid_whose_comm_is_not_java() {
|
||||
pgrep() { printf '4242\n'; }
|
||||
ps() { printf 'perl\n'; }
|
||||
local found
|
||||
found="$(running_pid)"
|
||||
unset -f pgrep ps
|
||||
[ -z "$found" ] \
|
||||
|| fail "running_pid() counted pid 4242 whose comm is 'perl', not 'java' — denying known shell names does not exclude a non-shell wrapper such as ssh or perl (fleetd #593 CORRECTION 1): found=[$found]"
|
||||
}
|
||||
|
||||
# fleetd #593 CORRECTION 1, hole 1 — pgrep can list a pid that exits before the following
|
||||
# `ps -o comm=` lookup runs; on a gone pid `ps` prints nothing, so `comm` comes back empty. Under
|
||||
# the round-1 denylist an empty string matched none of the denied shell names, so the dead pid was
|
||||
# still counted — the exact false-positive shape the ticket exists to remove, just rarer. The
|
||||
# allowlist fixes this for free: an empty comm is not `java` either.
|
||||
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup() {
|
||||
pgrep() { printf '4242\n'; }
|
||||
ps() { :; } # a pid that no longer exists: the real `ps -p <gone>` prints nothing and this mirrors that
|
||||
local found
|
||||
found="$(running_pid)"
|
||||
unset -f pgrep ps
|
||||
[ -z "$found" ] \
|
||||
|| fail "running_pid() counted pid 4242 whose comm lookup came back empty (the pid had already exited before the lookup ran) — an empty comm must not pass the allowlist (fleetd #593 CORRECTION 1): found=[$found]"
|
||||
}
|
||||
|
||||
# The positive backstop for both stubbed tests above, and for
|
||||
# test_running_pid_finds_a_real_java_named_second_process on whatever platform that live fixture
|
||||
# does not itself carry comm=java: the allowlist must still ADMIT the one comm value the real
|
||||
# daemon actually has. Measured on the real, currently-running daemon on this Mac: `comm=java`.
|
||||
test_running_pid_counts_a_pid_whose_comm_is_java() {
|
||||
pgrep() { printf '4242\n'; }
|
||||
ps() { printf 'java\n'; }
|
||||
local found
|
||||
found="$(running_pid)"
|
||||
unset -f pgrep ps
|
||||
printf '%s\n' "$found" | grep -qxF '4242' \
|
||||
|| fail "running_pid() did not count pid 4242 whose comm is 'java' — the daemon's own name must pass the allowlist: found=[$found]"
|
||||
}
|
||||
|
||||
# fleetd #593 instance 3 — assert_single_daemon's refusal message used to tell the operator to
|
||||
# "Investigate with 'pgrep -f \"\$PATTERN\"'", which — typed by hand or over ssh — is precisely the
|
||||
# self-matching invocation instance 2 above fixes. A source-text check, the same technique
|
||||
# test_no_error_lines_message_gated_by_drain_state uses: this is prose inside a die() call, never
|
||||
# reached by sourcing (the SOURCED guard stops before the main flow, and this text only prints
|
||||
# from inside a call assert_single_daemon makes when it is already refusing).
|
||||
test_die_message_does_not_recommend_bare_pgrep_as_remediation() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" block bad
|
||||
block="$(grep -A6 -F 'racing supervisor produces' "$src" || true)"
|
||||
[ -n "$block" ] || fail "could not find the assert_single_daemon refusal message in redeploy-fleetd.sh"
|
||||
bad="$(printf '%s' "$block" | grep -F "Investigate with 'pgrep -f" || true)"
|
||||
[ -z "$bad" ] \
|
||||
|| fail "assert_single_daemon's die message still hands the operator a bare 'pgrep -f \"\$PATTERN\"' as remediation (fleetd #593) — that is exactly the self-matching invocation"
|
||||
printf '%s' "$block" | grep -qF 'fleetd #593' \
|
||||
|| fail "assert_single_daemon's die message does not say in words that a pattern can match the caller (fleetd #593)"
|
||||
}
|
||||
|
||||
# fleetd #511 — jar_id()'s no-argument default was unpinned by any test: nothing proved it reports
|
||||
# $JAR (the live path) rather than $JAR_STAGED. Both halves matter, so this pins both: the bare call
|
||||
# must hash the live jar, and an explicit path argument must hash THAT file, not fall back to $JAR.
|
||||
@@ -646,6 +770,137 @@ test_wait_for_daemon_exit_times_out_if_pid_never_clears() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
|
||||
}
|
||||
|
||||
# fleetd #603 — wait_for_new_pid is the dual of wait_for_daemon_exit above: it must not report
|
||||
# success while running_pid() still answers empty, and must report success the moment a pid
|
||||
# appears. Same counter-file idiom as test_wait_for_daemon_exit_returns_true_once_pid_clears above,
|
||||
# for the same reason (running_pid() runs inside a `$(...)` subshell on every call).
|
||||
test_wait_for_new_pid_returns_true_once_pid_appears() {
|
||||
local counter_file="$TMP/wait-new-pid-calls" final_calls
|
||||
printf '0' > "$counter_file"
|
||||
running_pid() {
|
||||
local n
|
||||
n="$(cat "$counter_file")"
|
||||
n=$((n + 1))
|
||||
printf '%s' "$n" > "$counter_file"
|
||||
if [ "$n" -lt 3 ]; then printf ''; else printf '4242'; fi
|
||||
}
|
||||
sleep() { :; }
|
||||
wait_for_new_pid 10 || fail "wait_for_new_pid did not report success once the pid appeared"
|
||||
final_calls="$(cat "$counter_file")"
|
||||
[ "$final_calls" -ge 3 ] || fail "wait_for_new_pid returned before actually re-checking running_pid"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
|
||||
}
|
||||
|
||||
test_wait_for_new_pid_times_out_if_pid_never_appears() {
|
||||
local rc=0
|
||||
running_pid() { printf ''; }
|
||||
sleep() { :; }
|
||||
wait_for_new_pid 3 || rc=$?
|
||||
[ "$rc" -ne 0 ] || fail "wait_for_new_pid reported success while the pid never appeared"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
|
||||
}
|
||||
|
||||
# fleetd #603 — await_daemon_started folds the pid-appeared check and the healthz check into one
|
||||
# decision (see its own comment in redeploy-fleetd.sh for why a source-text grep cannot tell the two
|
||||
# possible behaviors apart here). These two tests are the ticket's own acceptance criteria, run
|
||||
# together in this one suite invocation so neither can be satisfied by code that never fails at all:
|
||||
#
|
||||
# 1. a slow start must still succeed — running_pid mimics a process that does not appear until
|
||||
# well after the OLD, buggy 10-second budget, but does appear, and healthz answers.
|
||||
# 2. a genuine failure must still fail, and still print the log tail — running_pid and the health
|
||||
# check both report nothing at all, ever.
|
||||
test_await_daemon_started_slow_pid_then_healthy_succeeds() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
stub_die_recorder
|
||||
local counter_file="$TMP/await-slow-pid-calls" output
|
||||
printf '0' > "$counter_file"
|
||||
running_pid() {
|
||||
local n
|
||||
n="$(cat "$counter_file")"
|
||||
n=$((n + 1))
|
||||
printf '%s' "$n" > "$counter_file"
|
||||
# Stays empty well past the old 10-second budget, then appears — the exact shape #603 reports.
|
||||
if [ "$n" -lt 12 ]; then printf ''; else printf '4242'; fi
|
||||
}
|
||||
sleep() { :; }
|
||||
poll_health_body() { printf '{"status":"ok"}'; return 0; }
|
||||
# NOT `output="$(await_daemon_started ...)"`: that would run the call in a subshell, and
|
||||
# NEW_PID — a plain global assignment inside the function, by design (see its own comment) — would
|
||||
# die with that subshell instead of reaching this test's own shell. Redirect to a file instead, the
|
||||
# same hazard await_daemon_started's own comment warns die() itself is subject to.
|
||||
await_daemon_started 60 "" "http://ignored/healthz" "$TMP/await-slow-pid.out" \
|
||||
> "$TMP/await-slow-pid-output.log" 2>&1
|
||||
output="$(cat "$TMP/await-slow-pid-output.log")"
|
||||
[ "$DIED_CALLED" = 0 ] \
|
||||
|| fail "await_daemon_started must not die on a slow-but-real start: $DIED_MESSAGE"
|
||||
assert_equals "4242" "$NEW_PID" "await_daemon_started NEW_PID after a slow-but-real start"
|
||||
printf '%s' "$output" | grep -qF 'started, pid 4242' \
|
||||
|| fail "await_daemon_started did not report the pid once it finally appeared"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
test_await_daemon_started_never_appears_dies_with_log_tail() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
stub_die_recorder
|
||||
local out_file="$TMP/await-never-appears.out"
|
||||
printf 'boot line one\nboot line two\n' > "$out_file"
|
||||
running_pid() { printf ''; }
|
||||
sleep() { :; }
|
||||
poll_health_body() { return 1; }
|
||||
await_daemon_started 2 "" "http://127.0.0.1:1/healthz" "$out_file" > /dev/null 2>&1
|
||||
[ "$DIED_CALLED" = 1 ] \
|
||||
|| fail "await_daemon_started must die when the daemon never appears and never becomes healthy"
|
||||
printf '%s' "$DIED_MESSAGE" | grep -qF 'never answered' \
|
||||
|| fail "await_daemon_started die message does not say healthz never answered"
|
||||
printf '%s' "$DIED_MESSAGE" | grep -qF 'boot line two' \
|
||||
|| fail "await_daemon_started die message does not include the log tail"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
# fleetd #603 review — the path the fall-through actually exists for, and the one gap a lead
|
||||
# mutation found in the first version of this test file: running_pid() NEVER finds anything (as its
|
||||
# own doc comment says it eventually will, once the daemon stops being launched as a plain
|
||||
# `java -jar` its allowlist recognises), while /healthz answers anyway. Neither of the two tests
|
||||
# above drives this: the slow-pid test has the pid appear, so the `else` branch never runs, and the
|
||||
# never-appears test fails BOTH checks, so it dies either way and cannot tell which branch fired.
|
||||
# This must not die, must warn (so the operator is told the pid could not be identified), and must
|
||||
# leave NEW_PID empty — the honest "could not establish this" answer, never a guessed pid, which is
|
||||
# what the final result line's `${NEW_PID:-unknown}` fallback exists to print truthfully.
|
||||
#
|
||||
# Proof this actually pins the behavior, not just the source text (paste from a real run, not
|
||||
# claimed): reverting the `warn` below back to `die "no process appeared"` (the old fleetd #603
|
||||
# defect, reintroduced) turns this one test red —
|
||||
# FAIL: await_daemon_started must not die when the pid is never found but healthz answers
|
||||
# — and restoring `warn` turns the whole suite green again. Both halves observed, not asserted.
|
||||
test_await_daemon_started_pid_never_found_but_healthy_warns_and_survives() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
stub_die_recorder
|
||||
local out_file="$TMP/await-pid-never-found.out" output
|
||||
printf 'boot line\n' > "$out_file"
|
||||
running_pid() { printf ''; }
|
||||
sleep() { :; }
|
||||
poll_health_body() { printf '{"status":"ok"}'; return 0; }
|
||||
# NOT `output="$(await_daemon_started ...)"` — see the slow-pid test above for why that would
|
||||
# drop NEW_PID's assignment in a subshell instead of reaching this test's own shell.
|
||||
await_daemon_started 3 "" "http://ignored/healthz" "$out_file" \
|
||||
> "$TMP/await-pid-never-found-output.log" 2>&1
|
||||
output="$(cat "$TMP/await-pid-never-found-output.log")"
|
||||
[ "$DIED_CALLED" = 0 ] \
|
||||
|| fail "await_daemon_started must not die when the pid is never found but healthz answers: $DIED_MESSAGE"
|
||||
printf '%s' "$output" | grep -qF 'falling through to the health check' \
|
||||
|| fail "await_daemon_started did not warn that the pid could not be identified"
|
||||
assert_equals "" "$NEW_PID" \
|
||||
"await_daemon_started NEW_PID when the pid is never found but healthz answers — must stay empty, never a guessed pid"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
test_await_daemon_started_call_site_present() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" call_line
|
||||
call_line="$(grep -Fn 'await_daemon_started "$HEALTH_WAIT" "$OLD_PID" "$HEALTH" "$OUT"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
[ -n "$call_line" ] \
|
||||
|| fail "could not find the main flow's await_daemon_started call site in redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
# fleetd #521 — the swap step's guard, at two levels.
|
||||
#
|
||||
# The first two tests call the predicate should_swap() directly. They pin its logic, and that is all
|
||||
@@ -1226,9 +1481,11 @@ test_poll_health_body_returns_nonzero_when_unreachable() {
|
||||
|
||||
test_report_health_call_site_present() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" call_line
|
||||
call_line="$(grep -Fn 'report_health "$HEALTH_BODY" "$HEALTH_CODE" "$OUT" "$HEALTH_WAIT"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
# fleetd #603 — moved from a literal main-flow call into await_daemon_started (see its own
|
||||
# comment above for why); this now finds the call inside that function instead.
|
||||
call_line="$(grep -Fn 'report_health "$HEALTH_BODY" "$HEALTH_CODE" "$out_file" "$health_wait"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
[ -n "$call_line" ] \
|
||||
|| fail "could not find the main flow's report_health call site in redeploy-fleetd.sh"
|
||||
|| fail "could not find await_daemon_started's report_health call site in redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
# fleetd #555 item 8 — `HAD_OLD_PID=0; [ -n "$OLD_PID" ] && HAD_OLD_PID=1`. Measured safe under
|
||||
@@ -1894,6 +2151,12 @@ test_require_drivable_supervisor_accepts_known_kinds
|
||||
test_count_daemon_pids
|
||||
test_assert_single_daemon_accepts_one_pid
|
||||
test_assert_single_daemon_rejects_two_pids
|
||||
test_running_pid_excludes_self_matching_wrapper_shell
|
||||
test_running_pid_finds_a_real_java_named_second_process
|
||||
test_running_pid_drops_a_pid_whose_comm_is_not_java
|
||||
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup
|
||||
test_running_pid_counts_a_pid_whose_comm_is_java
|
||||
test_die_message_does_not_recommend_bare_pgrep_as_remediation
|
||||
test_jar_id_defaults_to_live_and_reports_explicit_path
|
||||
test_hash256_computes_a_real_sha256
|
||||
test_jar_id_reports_absent_for_missing_file
|
||||
@@ -1911,6 +2174,12 @@ test_require_no_build_jar_dies_when_absent
|
||||
test_require_no_build_jar_accepts_present_jar
|
||||
test_wait_for_daemon_exit_returns_true_once_pid_clears
|
||||
test_wait_for_daemon_exit_times_out_if_pid_never_clears
|
||||
test_wait_for_new_pid_returns_true_once_pid_appears
|
||||
test_wait_for_new_pid_times_out_if_pid_never_appears
|
||||
test_await_daemon_started_slow_pid_then_healthy_succeeds
|
||||
test_await_daemon_started_never_appears_dies_with_log_tail
|
||||
test_await_daemon_started_pid_never_found_but_healthy_warns_and_survives
|
||||
test_await_daemon_started_call_site_present
|
||||
test_swap_ordered_after_wait_and_before_start
|
||||
test_drain_gate_abort_message_says_no_no_build
|
||||
test_drain_gate_refusal_build_ran_staged_present
|
||||
|
||||
Reference in New Issue
Block a user