Compare commits
20 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 26f198675c | |||
| 72f46d7c0c | |||
| 640f4d5f23 | |||
| 6edeb70bc4 | |||
| bc49d87cb8 | |||
| 14b169c410 | |||
| 2b52324d9a | |||
| b6147a39f6 | |||
| dbfc34cb6d | |||
| a20cb96730 | |||
| cda1a6a917 | |||
| 608e4496be | |||
| fa61dc587c | |||
| 526e3b7459 | |||
| b6b006c651 | |||
| 63eec8a0da | |||
| cbe872b538 | |||
| 7d9a807243 | |||
| 8915e40c7d | |||
| 3f7bc3815e |
@@ -167,19 +167,59 @@ fails.
|
||||
- **There is no terminal or session parameter, on purpose.** The pane is always your own, resolved
|
||||
from your connection, so you can only ever roll yourself.
|
||||
- **`operatorConfirmed` is your report of what a human told you.** Do not pass `true` because you
|
||||
are confident. Ask, wait for the answer, then pass what they said. `requireOperatorConfirm`
|
||||
defaults to `true` and this is the only thing standing between a judgement call and a wiped
|
||||
session.
|
||||
are confident. Ask, wait for the answer, then pass what they said.
|
||||
|
||||
- **Whether you must ask at all depends on `leadRollover.requireOperatorConfirm`. Check it; do not
|
||||
assume.** The default is `true` (`FleetConfig.java:1426`), and then `confirm` refuses unless you
|
||||
also pass `operatorConfirmed: true`. **This host set it to `false` on 2026-09-22**, on the
|
||||
operator's explicit grant, because they do not want to approve routine context rolls. Where it is
|
||||
`false`, the three handover-file checks are the whole gate: the file must exist, be fresher than
|
||||
`maxDocAgeSeconds`, and have been modified after the open request.
|
||||
|
||||
Read the live value rather than trusting this line:
|
||||
|
||||
```bash
|
||||
grep -A1 'requireOperatorConfirm' fleetd/fleetd.yaml
|
||||
```
|
||||
|
||||
No match means the key is unset, so the default `true` applies and you must ask. The key is
|
||||
**deferred, not hot** — it is read once at boot, so an edit does nothing until the daemon is
|
||||
redeployed.
|
||||
|
||||
**Until fleetd #621 merges, the nudge text will tell you to ask the operator even where the
|
||||
daemon no longer requires it.** `LeadHeartbeatLoop.contextNotice()` hardcodes "ask the operator"
|
||||
and takes no config, so it cannot know. Trust the config value over the nudge text. Once #621 is
|
||||
merged and deployed, the nudge matches the config and this warning can be deleted.
|
||||
|
||||
- **The roll can still refuse after `confirm` returns**, and by then there is no caller to tell.
|
||||
Those outcomes are logged only, as `lead-rollover:` lines in the daemon log.
|
||||
- **The bootstrap prompt has never yet landed, and the fix is unproven (fleetd #489).** The first
|
||||
real rollover, on 2026-09-12, joined `/clear` and the bootstrap text into one line and Claude Code
|
||||
refused it as `Unknown command: /clearFresh`. The pane was never cleared and no context was lost,
|
||||
so the failure was safe — the roll simply did nothing. PR #490 fixed the cause and is deployed,
|
||||
but no roll has bootstrapped a fresh session end to end yet. **Assume it may still fail, and tell
|
||||
the operator so before you confirm.** The recovery is the same either way: the file is already
|
||||
written, so the operator starts a session and points it at the file. That is why you write the
|
||||
file before you confirm, and never the other way round.
|
||||
- **The bootstrap prompt works end to end. Measured 2026-09-22.** This used to say the fix was
|
||||
unproven (fleetd #489) and told you to expect a failure. That is no longer true. The daemon log
|
||||
now holds four `lead-rollover: rolled` lines, and three of them ran on 2026-09-22 at 10:01:43,
|
||||
10:38:28 and 11:15:47. Each one cleared the old lead and started a fresh session against the
|
||||
handover file, with the configured `bootstrapText` arriving as its first message. No context was
|
||||
lost. The old `Unknown command: /clearFresh` failure from 2026-09-12 does not appear in the log
|
||||
at all. Re-measure both numbers with:
|
||||
|
||||
```bash
|
||||
grep -c "lead-rollover: rolled" fleetd/fleetd.out # successful rolls
|
||||
grep -c "lead-rollover:" fleetd/fleetd.out # positive control: must be larger
|
||||
grep -c "Unknown command" fleetd/fleetd.out # the old failure: expect 0
|
||||
```
|
||||
|
||||
Run the control line too. A broken pattern returns a clean `0` that reads exactly like good news.
|
||||
If the first number stops growing across rolls, or `Unknown command` returns anything above 0,
|
||||
the bootstrap has regressed and this paragraph is stale again.
|
||||
|
||||
**You still write the file before you confirm, and never the other way round.** That order is not
|
||||
about the bootstrap being unreliable. It is what the daemon checks: the handover file must have
|
||||
been modified *after* the open request, or `confirm` refuses it as stale.
|
||||
|
||||
- **One warning in the log is normal and is not a failure.** Every one of the three rolls above also
|
||||
logged `/clear on term_… was never observed as WORKING after 8 consecutive IDLE/DONE polls —
|
||||
releasing rather than wedging the roll`. The daemon could not see the pane go WORKING after
|
||||
`/clear`, so it released instead of hanging. The roll then succeeded anyway. That is the safe
|
||||
branch behaving correctly. Do not report it as a broken roll.
|
||||
|
||||
## Writing style
|
||||
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
|
||||
/**
|
||||
* fleetd #612 Unit A: everything {@code Fleetd.main} has ready once config is loaded, reported
|
||||
* and validated — the exact point {@code main} used to keep going straight into socket and broker
|
||||
* work. Building this (and handing it to {@link FleetdAssembly#assembleAndStart}) is the seam a
|
||||
* test now has to drive the real boot composition without being {@code main} itself.
|
||||
*
|
||||
* @param cfg the boot-time {@link FleetConfig} snapshot every one-time wiring decision reads —
|
||||
* see {@code Fleetd.main}'s own comment on why this must never be swapped for a live
|
||||
* reference once loaded
|
||||
* @param config the live {@link ConfigRef} the hot-reloadable paths read per use
|
||||
* @param guard the same {@link SubscriptionGuard} {@code Fleetd.main} already used to assert the
|
||||
* launching environment is clean, reused rather than rebuilt so the assembled
|
||||
* launchers see the identical instance {@code main} already validated with
|
||||
*/
|
||||
record AssemblyInputs(FleetConfig cfg, ConfigRef config, SubscriptionGuard guard) {
|
||||
}
|
||||
@@ -42,6 +42,7 @@ import dev.ltms.fleet.mcp.LsofProcessCwdLookup;
|
||||
import dev.ltms.fleet.msg.AmqpReplyInbox;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.LeadChannel;
|
||||
import dev.ltms.fleet.msg.LeadChannelHandle;
|
||||
import dev.ltms.fleet.msg.LeadCoordLoop;
|
||||
import dev.ltms.fleet.msg.LeadMailbox;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
@@ -177,613 +178,18 @@ public final class Fleetd {
|
||||
// and the Fleetd-startup tests actually pin — see FleetConfig#validateAll's javadoc for
|
||||
// why a name-by-name list here would have the same defect it replaces.
|
||||
cfg.validateAll();
|
||||
// fleetd #613: validateAll() (validateMembers() inside it) only refuses a slot that names a
|
||||
// bad role or profile — it says nothing about a role that has NO pool or NO charter at all,
|
||||
// because both are legitimate ("unconstrained") states, not errors. Report them here, right
|
||||
// after validation passes, so an operator sees the gap once per restart instead of finding
|
||||
// it later in a roster row (see reportRoleFallbackGaps' javadoc for the measured cause).
|
||||
reportRoleFallbackGaps(cfg);
|
||||
// fleetd #469, follow-up to #464: validateAll() (and validateCharters() inside it) only
|
||||
// checks that a charter's KEY is a role wire name and its text is non-blank — it never
|
||||
// looks at what the text actually names. This is the separate check that does: it asks
|
||||
// dev.ltms.fleet.mcp.FleetTool (the canonical registered-tool set) whether every fleet_*/
|
||||
// bridge_* token a charter names is a tool this server actually registers. It cannot live
|
||||
// inside FleetConfig#validateCharters() — config loads before the MCP server exists, and
|
||||
// must not gain a dependency on the mcp package — so it runs here instead, at the one seam
|
||||
// that already holds both a loaded FleetConfig and the mcp package, before anything below
|
||||
// opens a socket or spawns a member. fleetd #474: the same check is also wired into `config`
|
||||
// above as ConfigRef's extraValidation, so a reload refuses what this line refuses at startup.
|
||||
assertChartersNameOnlyRegisteredTools(cfg);
|
||||
|
||||
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
|
||||
? Path.of(cfg.herdrSocket())
|
||||
: UnixSocketHerdrClient.defaultSocketPath();
|
||||
|
||||
UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper());
|
||||
UnixSocketHerdrClient memberHerdr = cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank()
|
||||
? UnixSocketHerdrClient.connect(Path.of(cfg.memberHerdrSocket()), new com.fasterxml.jackson.databind.ObjectMapper())
|
||||
: herdr;
|
||||
AtomicReference<Supplier<Map<String, String>>> leadsRef = new AtomicReference<>(Map::of);
|
||||
HerdrRouter router = new HerdrRouter(herdr, memberHerdr,
|
||||
target -> leadsRef.get().get().containsKey(target));
|
||||
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
|
||||
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
|
||||
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
|
||||
Map<String, FleetConfig.Profile> claudeProfiles = new LinkedHashMap<>();
|
||||
Map<String, FleetConfig.Profile> opencodeProfiles = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, w) -> {
|
||||
if (w.isOpenCode()) {
|
||||
opencodeProfiles.put(name, w);
|
||||
} else {
|
||||
claudeProfiles.put(name, w);
|
||||
}
|
||||
});
|
||||
List<HerdrPeerLauncher> adapters = new ArrayList<>();
|
||||
// fleetd #175: the daemon's real ExhaustionSink can only be built once `sessions` exists
|
||||
// (below), but `sessions` needs `workers`, which needs the adapters built right here — a
|
||||
// genuine cycle. Break it exactly like liveCountRef below: a forwarding sink built now,
|
||||
// pointed at the real one once it exists.
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
// fleetd #234, round 4: routed through the shared ExhaustionSink.forwardingTo factory
|
||||
// rather than written inline here — not because a lambda at this call site is unsafe
|
||||
// anymore (it is not: the 3-arg overload is now the interface's single abstract method, so
|
||||
// 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).
|
||||
// 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(claudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
|
||||
claudeProfiles, cfg, config));
|
||||
}
|
||||
if (!opencodeProfiles.isEmpty()) {
|
||||
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
|
||||
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
|
||||
// The cooldown is deferred (see FleetConfig#quarantineCooldownSeconds): it is read once
|
||||
// here, at startup, and a config reload only changes it for a daemon restart.
|
||||
// fleetd #466: escalating, not flat — a credential that keeps reporting exhaustion (e.g. a
|
||||
// weekly subscription limit, which would otherwise be retried on every ~30-minute cooldown,
|
||||
// about 336 times across the week) backs off further each consecutive time, capped at
|
||||
// BackendQuarantine.DEFAULT_MAX_COOLDOWN_MULTIPLE x the base cooldown. See BackendQuarantine's
|
||||
// class doc for the mechanism, why this never fires on cooling-off (a separate, unescalated
|
||||
// mechanism — BackendOutagePolicy below), and the reset.
|
||||
BackendQuarantine quarantine = BackendQuarantine.withEscalation(System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
|
||||
// fleetd #201 Unit 5: one outage-cool-off tracker for the whole daemon, shared between the
|
||||
// launcher (checked at spawn, like `quarantine` above) and the backend-error sink wired in
|
||||
// below (written on a classified backend error). A SEPARATE, shorter-lived mechanism from
|
||||
// `quarantine` — see BackendOutagePolicy's class doc — never merged with it.
|
||||
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(System::nanoTime);
|
||||
PeerLauncher workers = new CompositePeerLauncher(
|
||||
adapters,
|
||||
cfg.effectiveDefaultProfile(),
|
||||
config,
|
||||
profileName -> liveCountRef.get().apply(profileName),
|
||||
quarantine,
|
||||
outagePolicy);
|
||||
// fleetd #422 follow-up: say which of the three model-gate states the daemon booted into —
|
||||
// no models: block at all, a block armed with nothing off, or a block with N off — the same
|
||||
// way exhaustedPatternCoverageLine/errorPatternCoverageLine report CB-578 stage A/fleetd
|
||||
// #201 Unit 5 coverage just below. Read from workers.modelGateState() (never a separate
|
||||
// config.get().models() here) so this line and fleet_profiles' modelGateArmed can never
|
||||
// disagree about what CompositePeerLauncher's spawn gate actually enforces.
|
||||
log.info("model gate (fleetd #422): {}", modelGateCoverageLine(workers.modelGateState()));
|
||||
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
|
||||
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
|
||||
// the first thing that actually talks to herdr, so without this wait a boot-order race
|
||||
// would crash the daemon into a restart loop. Wait, then degrade rather than die: serving
|
||||
// with /healthz reporting "degraded" is strictly more useful than exiting.
|
||||
HerdrAwaitOutcome herdrOutcome = awaitHerdr(herdr, System::nanoTime, Fleetd::sleepHerdrPoll);
|
||||
boolean herdrUp = logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
|
||||
if (herdrUp) {
|
||||
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
|
||||
// with the previous process — reap those leaked orphans now, before we start serving.
|
||||
workers.reapOrphanWorkers();
|
||||
}
|
||||
|
||||
// CB-301: authoritative session registry + lifecycle FSM on top of ClaudeCodeLauncher.
|
||||
// CB-301-ext: worktree provisioning seam, optionally rooted at a configured directory.
|
||||
// CB-303 part 2: context cap is opt-in and disabled (0) when absent/null.
|
||||
int contextCap = 0;
|
||||
if (cfg.lifecycle() != null && cfg.lifecycle().contextCap() != null
|
||||
&& cfg.lifecycle().contextCap() > 0) {
|
||||
contextCap = cfg.lifecycle().contextCap();
|
||||
}
|
||||
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
|
||||
SessionManager sessions = new SessionManager(workers,
|
||||
new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup(), cfg.memberSkills()),
|
||||
System::nanoTime, contextCap, clearAfterTurn);
|
||||
liveCountRef.set(profileName -> liveSessionCount(sessions.roster(), profileName));
|
||||
|
||||
// Idle-sleep guard: hold an OS-level assertion against idle sleep while at least one
|
||||
// member is live, so an unattended host does not idle-sleep out from under a member's
|
||||
// long turn (see FleetConfig.IdleSleepGuard / dev.ltms.fleet.power.IdleSleepGuard for the
|
||||
// measurement that motivated this). Opt-out via idleSleepGuard.enabled: false; on by
|
||||
// default. Hangs off SessionManager's own onAcquire/onRelease hooks (CB-520/CB-516,
|
||||
// previously wired only to the reply inbox) and SessionManager#size() — the exact registry
|
||||
// fleet_list's live/capacity numbers are themselves computed from — rather than tracking
|
||||
// members a second way. No-op (never constructed) off macOS or when idleSleepGuard.enabled
|
||||
// is explicitly false; the mechanism itself is additionally a no-op if 'caffeinate' cannot
|
||||
// be started, so this can never fail a spawn, a release, or startup.
|
||||
boolean idleSleepGuardEnabled = cfg.idleSleepGuard() == null || cfg.idleSleepGuard().isEnabled();
|
||||
final IdleSleepGuard idleSleepGuard;
|
||||
if (idleSleepGuardEnabled) {
|
||||
idleSleepGuard = new IdleSleepGuard(new CaffeinateSleepAssertionMechanism(), sessions::size);
|
||||
sessions.onAcquire(_ -> idleSleepGuard.recheck());
|
||||
sessions.onRelease(_ -> idleSleepGuard.recheck());
|
||||
} else {
|
||||
idleSleepGuard = null;
|
||||
}
|
||||
|
||||
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
|
||||
final SessionReaper reaper;
|
||||
if (cfg.lifecycle() != null
|
||||
&& cfg.lifecycle().idleTtlSeconds() != null
|
||||
&& cfg.lifecycle().idleTtlSeconds() > 0) {
|
||||
reaper = new SessionReaper(sessions, cfg.lifecycle().idleTtlSeconds());
|
||||
reaper.start();
|
||||
} else {
|
||||
reaper = null;
|
||||
}
|
||||
|
||||
// CB-530: every pane the config names as a lead, merged from `leaders:` and the legacy
|
||||
// singular pin. PrimaryRegistry below still tracks ONE terminal — it addresses the push
|
||||
// loop's nudges, which need a single destination — so it keeps the legacy pin.
|
||||
Map<String, String> leadTerminals = cfg.leaderTerminals();
|
||||
if (leadTerminals.size() > 1) {
|
||||
log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values());
|
||||
}
|
||||
// CB-531: on top of the legacy primary.terminal pin, discover leads by the tab labels the
|
||||
// operator writes. CB-557 moved the settings onto the lead they describe, so scanning is on
|
||||
// whenever a `fleet.leaders:` entry exists — with no leads configured the supplier is a
|
||||
// constant and never touches herdr, exactly as a missing `leadScan:` block used to behave.
|
||||
// CB-579: each lead now names its own exact `tab:` label, so one scanner discovers every
|
||||
// configured lead regardless of how differently their tabs are labelled — the old
|
||||
// single-shared-tabPrefix limitation (and its warning) is gone.
|
||||
final Supplier<Map<String, String>> leads;
|
||||
var leaders = cfg.fleet().leaders();
|
||||
if (!leaders.isEmpty()) {
|
||||
Map<String, String> tabToName = new LinkedHashMap<>();
|
||||
leaders.forEach((name, leader) -> {
|
||||
if (leader != null && leader.tab() != null && !leader.tab().isBlank()) {
|
||||
tabToName.put(leader.tab(), name);
|
||||
}
|
||||
});
|
||||
// The lead and the members now share ONE workspace (the operator asked for a single
|
||||
// "session" with many tabs), so no workspace can be excluded — the lead lives in the
|
||||
// members' space by design. A lead is told from a member by its exact tab label alone:
|
||||
// a lead carries its configured `tab`, a member its `worker: {profile} #{n}` template,
|
||||
// and the two never collide. (The scanner still supports an exclusion set for a split
|
||||
// layout; the fleet's policy here is simply not to use one.)
|
||||
//
|
||||
// One shared rescan cadence: still taken from the first entry, as before — it is an
|
||||
// operational cadence, not identity, so there is no correctness reason to give every
|
||||
// lead its own scanner.
|
||||
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
|
||||
// This must use the lead daemon: scanning member tabs would demote the lead to a worker.
|
||||
leads = new LeadTabScanner(herdr, tabToName, Set.of(),
|
||||
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), System::nanoTime);
|
||||
log.info("lead scan: tabs {} host a lead (rescan every {}s, shared fleet space)",
|
||||
tabToName.keySet(), scanIntervalSeconds);
|
||||
} else {
|
||||
leads = () -> leadTerminals;
|
||||
}
|
||||
leadsRef.set(leads);
|
||||
|
||||
// CB-558: start any declared lead that is not already running. After the scanner is built,
|
||||
// because both read the same tab labels and the ordering makes that dependency visible; and
|
||||
// only when herdr answered, because the launcher's whole safety property is that it can
|
||||
// count live leads first — it must never guess and risk a second orchestrator.
|
||||
if (herdrUp && !leaders.isEmpty()) {
|
||||
int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads();
|
||||
if (launched > 0) {
|
||||
log.info("lead auto-launch: {} lead(s) started", launched);
|
||||
}
|
||||
}
|
||||
|
||||
// CB-548: config-declared architect slots. Config supplies only the stable name → profile
|
||||
// map; the terminal → slot binding is owned by the registry and is empty at startup, so no
|
||||
// pane resolves to an architect until the later spawn lifecycle binds one. The registry is
|
||||
// what CallerResolver resolves against and what that lifecycle will read profiles from;
|
||||
// nothing here spawns a slot.
|
||||
// fleetd #424: MemberRegistry.live re-reads fleet.architects through `config` on every
|
||||
// reserve/requireSlotFor call, so a reload that removes or adds an architect slot governs
|
||||
// the next spawn with no restart — the frozen `new MemberRegistry(cfg.fleet())` this used
|
||||
// to be let a "revoked" slot keep granting new architect spawns forever.
|
||||
MemberRegistry members = MemberRegistry.live(() -> config.get().fleet());
|
||||
sessions.setMemberLifecycle(members);
|
||||
if (!members.slots().isEmpty()) {
|
||||
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
|
||||
+ "spawn lifecycle binds a live terminal to it)",
|
||||
members.slots().size(), members.slots().keySet());
|
||||
}
|
||||
|
||||
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
|
||||
// The blocking message endpoint (CB-104) is the producer; the poller is inert until then.
|
||||
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
// CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured
|
||||
// usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer.
|
||||
// fleetd #446: read live off `config` per lookup, cached by profile name — see
|
||||
// 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.
|
||||
// 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.
|
||||
Set<String> exhaustedConfiguredAtStartup = cfg.profiles().entrySet().stream()
|
||||
.filter(e -> e.getValue().hasExhaustedPattern())
|
||||
.map(Map.Entry::getKey)
|
||||
.collect(Collectors.toCollection(LinkedHashSet::new));
|
||||
log.info("backend-exhausted classification (CB-578 stage A): {}",
|
||||
exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedConfiguredAtStartup));
|
||||
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
|
||||
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
|
||||
// rather than handing it back as a real answer. Still compiled once at startup, keyed by
|
||||
// profile name — unlike exhaustedPattern above (fleetd #446), errorPattern was left DEFERRED
|
||||
// on purpose: the ticket that made exhaustedPattern hot scoped errorPattern/cooling-off out
|
||||
// explicitly. A profile with no configured errorPattern is simply absent here, so
|
||||
// CompletionResolver falls back to its built-in narrow {@code (?i)\bAPI Error\s*:}
|
||||
// compatibility pattern for that profile's targets (BackendErrorPatternLookup#legacy's
|
||||
// contract — see backendErrorPatterns below).
|
||||
Map<String, Pattern> errorPatternsByProfile = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, profile) -> {
|
||||
if (profile.hasErrorPattern()) {
|
||||
errorPatternsByProfile.put(name, Pattern.compile(profile.errorPattern()));
|
||||
}
|
||||
});
|
||||
// fleetd #248: extracted to a static factory (see backendErrorPatternLookup below) so a
|
||||
// test can prove main() actually PASSES this into CompletionResolver, not only that the
|
||||
// lookup itself behaves correctly — the exact gap fleetd #248 exists to close.
|
||||
BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,
|
||||
errorPatternsByProfile);
|
||||
log.info("backend-error classification (fleetd #201 Unit 5): {}",
|
||||
errorPatternCoverageLine(cfg.profiles().keySet(), errorPatternsByProfile.keySet()));
|
||||
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
|
||||
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
|
||||
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
|
||||
// the profile config live off `config`, so a credentialId edit is hot: no restart needed.
|
||||
// fleetd #446 criterion 3: the backend text that triggered the most recent quarantine of
|
||||
// each credential, so fleet_profiles can report WHY a limit was hit, not only that it was.
|
||||
// Keyed by credential id — the same key BackendQuarantine's own remainingSeconds uses —
|
||||
// and written at the one call site that actually quarantines (inside exhaustionSink below),
|
||||
// so a reason can never be reported for a quarantine that never happened. Bounded the same
|
||||
// way BackendQuarantine's own internal map is documented to be: by the number of distinct
|
||||
// credentials ever exhausted, not by how often the config is edited.
|
||||
Map<String, String> quarantineReasonByCredential = new ConcurrentHashMap<>();
|
||||
// fleetd #446 follow-up (round 3): extracted into exhaustionSink(...) below — see that
|
||||
// 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.
|
||||
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
|
||||
// now that `sessions` exists to resolve target -> session -> profile.
|
||||
// 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
|
||||
// holder set once `pushLoop` exists, read lazily from inside the lambda built here.
|
||||
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>();
|
||||
// fleetd #248: extracted to a static factory (see backendErrorSink below), public rather
|
||||
// than package-private like the other two factories here, so
|
||||
// dev.ltms.fleet.inject.BackendOutageFlowTest can exercise the REAL production sink
|
||||
// directly instead of a hand-mirrored copy of this lambda — that copy was precisely the
|
||||
// gap fleetd #248 exists to close (see that test's class doc for the history).
|
||||
BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),
|
||||
outagePolicy, pushLoopRef::get);
|
||||
AgentControl agents = router.memberAgents();
|
||||
// Both fleetd#201 Unit 5 (backend-error patterns + sink) and fleetd#241 (the worktree/branch
|
||||
// lookup the fallback report names) land on this one call. The full constructor takes both,
|
||||
// so neither feature is dropped; nowNanos must be passed explicitly to reach it.
|
||||
//
|
||||
// fleetd #248: every argument built specifically for this call (backendErrorPatterns and
|
||||
// backendErrorSink above, and the worktree/branch lookup right here) now comes from a
|
||||
// static factory tested on its own; FleetdCompletionResolverWiringTest source-asserts that
|
||||
// THIS call actually passes them, which is the coverage that was missing before.
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns,
|
||||
exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,
|
||||
worktreeBranchLookup(sessions::roster));
|
||||
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
||||
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
|
||||
MemberPresence presence = sessions.asPresence();
|
||||
// fleetd #561: extracted to a static factory (see turnListener below) — the anonymous class
|
||||
// this replaced had two bare, unguarded statements per callback, and nothing enforced that
|
||||
// the completion half went first beyond call order in the source.
|
||||
TurnListener turnListener = turnListener(completion, sessions);
|
||||
Predicate<String> deliverable = deliverableTo(presence, leads);
|
||||
// fleetd #556: registration is wired directly to `completion`, not folded into the
|
||||
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
|
||||
// listener) throwing, regardless of call order. See TurnRegistrar's javadoc.
|
||||
// fleetd #589 Group 3 (:505): extracted to turnRegistrar(...) — see FleetdTurnRegistrarWiringTest.
|
||||
Injector injector = new Injector(router, turnListener, deliverable,
|
||||
presence::forget, 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.
|
||||
// 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.
|
||||
// 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.
|
||||
String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null;
|
||||
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal);
|
||||
// CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs —
|
||||
// identity comes from `leaders:`/`leadScan:`, and reply nudges now follow the delegating
|
||||
// lead. Say so once at startup rather than leaving a redundant pin to look load-bearing.
|
||||
if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) {
|
||||
log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes "
|
||||
+ "from leaders:/leadScan:, and reply nudges follow the lead that delegated. "
|
||||
+ "It still works, and is still the fallback nudge destination when a restart "
|
||||
+ "has lost the delegation map. Its pushReminders/pushBackoffMs stay valid.");
|
||||
}
|
||||
// CB-307: active push-to-primary loop — nudge the primary when replies land without an
|
||||
// open fleet_send. Uses its own lightweight scheduled executor, separate from the injector.
|
||||
int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
|
||||
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
|
||||
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
||||
Thread.ofVirtual().name("bridge-push-").unstarted(r));
|
||||
// CB-502: the registry is built before the service and the push loop so send/reply outcomes
|
||||
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
|
||||
// CB-512: the push loop takes it too, so nudge outcomes (delivered|exhausted) are counted.
|
||||
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
|
||||
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
|
||||
pushScheduler, maxReminders, backoffMs, metrics);
|
||||
// fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda
|
||||
// above at the real push loop, now that it exists.
|
||||
pushLoopRef.set(pushLoop);
|
||||
// CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so
|
||||
// an upgraded daemon cannot silently start spending subscription on nudging an idle lead.
|
||||
// It has its own single-thread scheduler and holds its own scheduler shutdown via close().
|
||||
final LeadHeartbeatLoop heartbeat;
|
||||
var heartbeatScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
||||
Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r));
|
||||
if (cfg.leadHeartbeat() != null) {
|
||||
var hb = cfg.leadHeartbeat();
|
||||
// fleetd #609: own LeadContextGauge instance for the heartbeat loop — separate from the
|
||||
// one FleetMcp builds internally for fleet_list's context row. Each caches independently
|
||||
// (keyed by configDir+sessionId), so this costs at most one extra bounded tail read per
|
||||
// TTL window, never a shared-mutable-state hazard between the two callers.
|
||||
var leadContextGauge = new LeadContextGauge();
|
||||
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
|
||||
pushLoop, heartbeatScheduler, System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
|
||||
metrics,
|
||||
leadContextSource(leadContextGauge, router.leadAgents(), leads,
|
||||
leadConfigDirLookup(() -> config.get().profiles(), leaders)),
|
||||
Boolean.TRUE.equals(hb.contextHighNudge()));
|
||||
heartbeat.start();
|
||||
} else {
|
||||
heartbeat = null;
|
||||
heartbeatScheduler.shutdownNow();
|
||||
}
|
||||
// fleetd #480: lead rollover. Opt-in; absent `leadRollover:` this is never constructed, so
|
||||
// an upgraded daemon cannot silently acquire the ability to clear the lead's own pane.
|
||||
// Unlike heartbeat above, this has no recurring scheduler of its own — nothing but an
|
||||
// explicit confirm() call (wired to an MCP tool by a later ticket; nothing calls it yet)
|
||||
// that passes every gate can ever schedule a roll. It does not take primaryRegistry: the
|
||||
// lead terminal to roll comes from the caller of open()/confirm() (resolved by the MCP
|
||||
// layer from the connection, the same way auth/CallerResolver#resolve builds a
|
||||
// Principal.leader(...)), never from a single-slot lookup — see LeadRollover's class
|
||||
// javadoc, fleetd #480 correction 2.
|
||||
LeadRollover leadRollover = leadRollover(cfg, router.leadAgents(), config, leads);
|
||||
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
|
||||
pushLoop, metrics);
|
||||
|
||||
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
|
||||
final FleetHealthMonitor healthMonitor;
|
||||
var healthScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
||||
Thread.ofVirtual().name("bridge-health-").unstarted(r));
|
||||
if (cfg.health() != null && cfg.health().isEnabled()) {
|
||||
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
|
||||
// through the same idempotent target-wide operation CB-516 already uses on release.
|
||||
// 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(), healthFailTarget(messages));
|
||||
String coverage = FleetHealthMonitor.coverage(true,
|
||||
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
||||
if ("detection-only".equals(coverage)) {
|
||||
log.warn("fleet health: {} (no notification sink configured)", coverage);
|
||||
} else {
|
||||
log.info("fleet health: {}", coverage);
|
||||
}
|
||||
healthMonitor.start();
|
||||
} else {
|
||||
healthMonitor = null;
|
||||
healthScheduler.shutdownNow();
|
||||
}
|
||||
|
||||
// CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
|
||||
// release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
|
||||
sessions.onAcquire(replyInbox::own);
|
||||
// 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.
|
||||
// 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.
|
||||
// CB-185: a caller's pane can live on either daemon (a lead's on the lead daemon, a
|
||||
// member's on the member daemon) — search both, lead first. Collapses to one scan when
|
||||
// memberHerdrSocket is unset (herdr == memberHerdr).
|
||||
ConnectionIdentity identity = new ConnectionIdentity(
|
||||
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||
|
||||
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
|
||||
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
|
||||
final CallerResolver callers;
|
||||
if (cfg.auth().tokenMode()) {
|
||||
String token = System.getenv(cfg.auth().tokenEnv());
|
||||
if (token == null || token.isBlank()) {
|
||||
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
|
||||
+ " is unset or empty — export it before starting fleetd");
|
||||
}
|
||||
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members);
|
||||
log.info("auth: token mode (bearer required for non-worker callers, env {})",
|
||||
cfg.auth().tokenEnv());
|
||||
} else {
|
||||
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members);
|
||||
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
|
||||
}
|
||||
|
||||
// fleetd #297: named once and reused verbatim below for FleetApp's GET /profiles, rather than
|
||||
// built a second time — two independently-constructed sources reading the SAME BackendQuarantine
|
||||
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
|
||||
// closure in only one of the two places), exactly the shape #284 was.
|
||||
FleetMcp.QuarantineSource quarantineSource = quarantineSource(config, quarantine,
|
||||
liveExhaustedPatterns, quarantineReasonByCredential);
|
||||
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy);
|
||||
|
||||
FleetMcp.LoopHealthSource loopHealth = loopHealthSource(poller, reaper);
|
||||
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
|
||||
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
|
||||
capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
|
||||
healthCoverageSource(config),
|
||||
loopHealth,
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
|
||||
// fleetd #602 gauge-wiring: threads each lead's configured configDir into the
|
||||
// context gauge — see leadConfigDirSource's own doc for why this, not a hardcoded
|
||||
// null, is what fleet_list's context row now reads.
|
||||
leadConfigDirSource(() -> config.get().profiles(), leaders),
|
||||
// fleetd #361: the operator-declared peers this daemon's fleet_list should try to
|
||||
// reach. Read from the SAME snapshot leadMailbox itself opened from (cfg.coordinator()),
|
||||
// not the live config.get() — coordinator wiring is already boot-time-fixed (see
|
||||
// leadMailbox above), so peers follows the same rule rather than half hot-reloading.
|
||||
cfg.coordinator() == null ? List.of() : cfg.coordinator().peers(),
|
||||
// fleetd #480 Unit C: the executor behind fleet_handover — null whenever
|
||||
// leadRollover: is not configured (see the leadRollover local above).
|
||||
leadRollover);
|
||||
|
||||
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
|
||||
// coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is
|
||||
// created and no thread runs. It reads the SAME live lead supplier the injector's
|
||||
// deliverability gate does, so a lead found by the tab scan after startup is reachable
|
||||
// without a restart.
|
||||
final LeadCoordLoop leadCoordLoop;
|
||||
final ScheduledExecutorService leadCoordSchedulerRef;
|
||||
if (leadMailbox != null) {
|
||||
var leadCoordScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
||||
Thread.ofVirtual().name("bridge-leadcoord-").unstarted(r));
|
||||
leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler,
|
||||
LEAD_COORD_INTERVAL_MS);
|
||||
leadCoordLoop.start();
|
||||
leadCoordSchedulerRef = leadCoordScheduler;
|
||||
} else {
|
||||
leadCoordLoop = null;
|
||||
leadCoordSchedulerRef = null;
|
||||
}
|
||||
|
||||
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an
|
||||
// upgraded daemon behaves exactly as before — the file is read once at boot and never again.
|
||||
final ConfigWatcher configWatcher;
|
||||
if (cfg.configReload() != null && cfg.configReload().isEnabled()) {
|
||||
configWatcher = new ConfigWatcher(config, cfg.configReload().intervalSeconds());
|
||||
configWatcher.start();
|
||||
} else {
|
||||
configWatcher = null;
|
||||
}
|
||||
|
||||
// CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still
|
||||
// open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close herdr
|
||||
// last. This replaces the earlier independent hooks that could race and close herdr early.
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
|
||||
sessions.close(cfg.lifecycle() != null ? cfg.lifecycle().drainTimeoutSeconds() : null);
|
||||
poller.stop();
|
||||
messages.close();
|
||||
pushLoop.close();
|
||||
if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler
|
||||
if (leadCoordLoop != null) leadCoordLoop.close(); // CB-637: stop delivering peer-lead messages
|
||||
if (leadCoordSchedulerRef != null) leadCoordSchedulerRef.shutdownNow();
|
||||
if (healthMonitor != null) healthMonitor.stop();
|
||||
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
|
||||
mcp.close();
|
||||
if (reaper != null) reaper.stop();
|
||||
// Idle-sleep guard: release unconditionally, even though sessions.close() above already
|
||||
// drained every session (and each release already drove the live count to 0, which
|
||||
// releases the guard's assertion on its own) — this is the backstop for a drain that was
|
||||
// itself interrupted or threw, so no caffeinate child ever outlives the daemon.
|
||||
if (idleSleepGuard != null) idleSleepGuard.close();
|
||||
// Release the broker connection last among message resources (no-op for the in-memory inbox).
|
||||
if (replyInbox instanceof AutoCloseable closeable) {
|
||||
try {
|
||||
closeable.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("reply inbox close: {}", e.toString());
|
||||
}
|
||||
}
|
||||
// CB-637: the coordination connection goes with it — after the loop that reads it has
|
||||
// stopped, so no tick can be mid-ack against a closed channel.
|
||||
if (leadMailbox != null) {
|
||||
try {
|
||||
leadMailbox.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("lead mailbox close: {}", e.toString());
|
||||
}
|
||||
}
|
||||
router.close();
|
||||
}));
|
||||
|
||||
// CB-185: give FleetApp both daemons — /healthz must require both to answer and
|
||||
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
|
||||
// fleetd #111: live (re-read-per-request) memberCredentials view for GET /member-credentials —
|
||||
// same hot-reload shape as the memberCredentials supplier passed to ClaudeCodeLauncher above.
|
||||
// fleetd #297: quarantineSource/outageSource are the SAME instances passed to FleetMcp above —
|
||||
// GET /profiles must report the identical quarantine/cool-off facts as fleet_profiles.
|
||||
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
|
||||
callers, metrics, deliverable,
|
||||
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
|
||||
quarantineSource, outageSource, loopHealth).build();
|
||||
app.start(cfg.bind().host(), cfg.bind().port());
|
||||
log.info("fleetd listening on {}:{}, herdr socket {}",
|
||||
cfg.bind().host(), cfg.bind().port(), socket);
|
||||
// fleetd #612 A-gaps (gap 2): everything from here on — including the two post-validation
|
||||
// reports that used to run inline right here (reportRoleFallbackGaps,
|
||||
// assertChartersNameOnlyRegisteredTools) — now lives in FleetdAssembly.assembleAndStart,
|
||||
// built against a real ResourcePorts. Unit A originally moved only the socket/broker/HTTP
|
||||
// composition and left those two calls here, between validateAll() and the assembly call —
|
||||
// outside the boundary FleetdAssemblyLifecycleTest drives, so deleting either call
|
||||
// compiled clean and left the whole suite green. Moving the boundary to start immediately
|
||||
// after validateAll() (this line) puts both back under test, in the same relative order,
|
||||
// before either one does any I/O — see FleetdAssembly's javadoc for the full boot-order
|
||||
// contract this preserves exactly.
|
||||
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ResourcePorts.system());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1809,10 +1215,18 @@ public final class Fleetd {
|
||||
ReplyInbox open(String uri, int prefetch);
|
||||
}
|
||||
|
||||
/** Injection seam for {@link #openLeadMailbox}: production binds {@link LeadMailbox#open}. */
|
||||
/**
|
||||
* Injection seam for {@link #openLeadMailbox}: production binds {@link LeadMailbox#open}.
|
||||
*
|
||||
* <p>fleetd #612 A-gaps (gap 1): returns {@link LeadChannelHandle}, not the concrete {@link
|
||||
* LeadMailbox}, so a test can supply a fake closeable channel instead of a real broker
|
||||
* connection — see {@link LeadChannelHandle}'s own javadoc for why the narrower type existed
|
||||
* and what widening it to add {@code close()} costs (nothing: {@code LeadMailbox} already
|
||||
* implements it).
|
||||
*/
|
||||
@FunctionalInterface
|
||||
interface LeadMailboxOpener {
|
||||
LeadMailbox open(String uri, String selfCoordId, int prefetch);
|
||||
LeadChannelHandle open(String uri, String selfCoordId, int prefetch);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1835,7 +1249,7 @@ public final class Fleetd {
|
||||
* fleet still works exactly as it did before this feature existed.</li>
|
||||
* </ul>
|
||||
*/
|
||||
static LeadMailbox openLeadMailbox(FleetConfig.Coordinator coordinator, Map<String, String> env,
|
||||
static LeadChannelHandle openLeadMailbox(FleetConfig.Coordinator coordinator, Map<String, String> env,
|
||||
LeadMailboxOpener opener) {
|
||||
if (coordinator == null) {
|
||||
return null; // opt-in: nothing configured, nothing to say
|
||||
@@ -1855,7 +1269,7 @@ public final class Fleetd {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
LeadMailbox mailbox = opener.open(uri, coordinator.selfId(), coordinator.prefetchOrDefault());
|
||||
LeadChannelHandle mailbox = opener.open(uri, coordinator.selfId(), coordinator.prefetchOrDefault());
|
||||
log.info("lead coordination: ON as coord-id {} (prefetch={})",
|
||||
coordinator.selfId(), coordinator.prefetchOrDefault());
|
||||
return mailbox;
|
||||
@@ -2235,16 +1649,20 @@ public final class Fleetd {
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #474: the one place both the startup call (right after {@code cfg.validateAll()} in
|
||||
* {@link #main}) and the reload call (wired into {@code config}'s {@code extraValidation} above,
|
||||
* via a method reference to this method) go through, so the two can never drift into checking
|
||||
* different things. Extracted only to give {@link ConfigRef}'s {@code Consumer<FleetConfig>}
|
||||
* hook a {@code FleetConfig -> void} shape to bind to — {@link CharterToolSurface} itself still
|
||||
* takes the raw charter map and knows nothing about {@code ConfigRef} or {@code Fleetd}.
|
||||
* fleetd #474: the one place both the startup call (fleetd #612 A-gaps gap 2: right after
|
||||
* {@code cfg.validateAll()} in {@link FleetdAssembly#assembleAndStart}, immediately after
|
||||
* {@code main} hands off to it — see that method's javadoc) and the reload call (wired into
|
||||
* {@code config}'s {@code extraValidation} above, via a method reference to this method) go
|
||||
* through, so the two can never drift into checking different things. Extracted only to give
|
||||
* {@link ConfigRef}'s {@code Consumer<FleetConfig>} hook a {@code FleetConfig -> void} shape to
|
||||
* bind to — {@link CharterToolSurface} itself still takes the raw charter map and knows nothing
|
||||
* about {@code ConfigRef} or {@code Fleetd}.
|
||||
*
|
||||
* <p>Package-private so a test can call it directly the same way the other startup-report
|
||||
* helpers above are tested, without needing to drive {@link #main} for a unit-level check;
|
||||
* {@code FleetdStartupValidationTest} proves the startup call site, and {@code
|
||||
* {@code FleetdStartupValidationTest} proves the startup call site by driving {@code main}
|
||||
* itself end to end (the throw still happens before {@code main} reaches any real socket or
|
||||
* broker work, since the assembly runs this before either), and {@code
|
||||
* FleetdConfigRefCharterToolSurfaceWiringTest} — by constructing {@code ConfigRef} with this
|
||||
* exact method reference, the same way {@code main} does above — proves the reload call site.
|
||||
* {@code dev.ltms.fleet.config.ConfigRefTest} pins the same reload behaviour too, through an
|
||||
@@ -2286,13 +1704,16 @@ public final class Fleetd {
|
||||
record HerdrAwaitOutcome(HerdrWaitResult result, long elapsedNanos) {}
|
||||
|
||||
/**
|
||||
* The real per-poll wait {@link #main} passes to {@link #awaitHerdr}: sleep
|
||||
* {@link #HERDR_WAIT_POLL_MILLIS}, and on interruption re-set the thread's interrupt flag
|
||||
* The real per-poll wait {@link FleetdAssembly#assembleAndStart} passes to {@link #awaitHerdr}:
|
||||
* sleep {@link #HERDR_WAIT_POLL_MILLIS}, and on interruption re-set the thread's interrupt flag
|
||||
* rather than throwing — {@link #awaitHerdr} detects an interruption by checking {@link
|
||||
* Thread#isInterrupted()} right after this returns, so a poller that swallowed the flag
|
||||
* instead of restoring it would make that check silently miss the interruption.
|
||||
*
|
||||
* <p>Package-private (fleetd #612 Unit A), not {@code private}: {@code FleetdAssembly}, a
|
||||
* different class in this package, needs a {@code Runnable} reference to this exact method.
|
||||
*/
|
||||
private static void sleepHerdrPoll() {
|
||||
static void sleepHerdrPoll() {
|
||||
try {
|
||||
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
|
||||
} catch (InterruptedException ie) {
|
||||
|
||||
@@ -0,0 +1,536 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.ConfigWatcher;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.health.FleetHealthMonitor;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.herdr.LeadTabScanner;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.herdr.UnixSocketHerdrClient;
|
||||
import dev.ltms.fleet.inject.BackendErrorPatternLookup;
|
||||
import dev.ltms.fleet.inject.BackendErrorSink;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
|
||||
import dev.ltms.fleet.inject.MemberPresence;
|
||||
import dev.ltms.fleet.inject.StatusPoller;
|
||||
import dev.ltms.fleet.inject.TurnListener;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.lead.LeadLauncher;
|
||||
import dev.ltms.fleet.lead.LeadRollover;
|
||||
import dev.ltms.fleet.mcp.ConnectionIdentity;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.mcp.LsofPeerPidLookup;
|
||||
import dev.ltms.fleet.mcp.LsofProcessCwdLookup;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.member.CompositePeerLauncher;
|
||||
import dev.ltms.fleet.member.HerdrPeerLauncher;
|
||||
import dev.ltms.fleet.member.MemberCredentialPolicyView;
|
||||
import dev.ltms.fleet.metrics.FleetMetrics;
|
||||
import dev.ltms.fleet.metrics.Metrics;
|
||||
import dev.ltms.fleet.msg.LeadChannelHandle;
|
||||
import dev.ltms.fleet.msg.LeadCoordLoop;
|
||||
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.msg.ReplyPushLoop;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.power.CaffeinateSleepAssertionMechanism;
|
||||
import dev.ltms.fleet.power.IdleSleepGuard;
|
||||
import dev.ltms.fleet.rest.FleetApp;
|
||||
import dev.ltms.fleet.session.GitWorktrees;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.SessionReaper;
|
||||
import io.javalin.Javalin;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
* fleetd #612 Unit A: the real boot assembly, extracted out of {@code Fleetd.main} so a test can
|
||||
* drive it directly. {@link #assembleAndStart} is <em>the same statements {@code main} used to run
|
||||
* inline</em>, in the same order, against a real {@link ResourcePorts} in production and a fake one
|
||||
* in a test — see {@code FleetdAssemblyLifecycleTest}. {@code Fleetd.main} keeps config loading and
|
||||
* {@code cfg.validateAll()}; everything from immediately after that call onward moved here —
|
||||
* including, since fleetd #612 A-gaps (gap 2), the two post-validation reports ({@code
|
||||
* reportRoleFallbackGaps}, {@code assertChartersNameOnlyRegisteredTools}) that Unit A originally
|
||||
* left behind in {@code main}. Those two calls do no I/O themselves, but leaving them outside this
|
||||
* method meant deleting either one compiled clean and left the whole suite green — nothing drove
|
||||
* {@code main} itself, so nothing could notice. They run first here, in the same relative order,
|
||||
* before the herdr socket or anything else that touches the outside world.
|
||||
*
|
||||
* <p><strong>Construction and start order is preserved exactly, on purpose.</strong> This is not
|
||||
* rebuilt into "construct everything, then start everything" — that would change boot timing. The
|
||||
* order recorded before any code moved (see the ticket and {@code FleetdAssemblyLifecycleTest}):
|
||||
* {@code SessionReaper} starts first (if {@code lifecycle.idleTtlSeconds} is configured), then
|
||||
* {@link StatusPoller}, then the optional {@link LeadHeartbeatLoop} and {@link FleetHealthMonitor},
|
||||
* then the optional {@link LeadCoordLoop} and {@link ConfigWatcher}, and the HTTP server starts
|
||||
* last of all. The close order (see {@link FleetdRuntime#close()}) is the mirror the original
|
||||
* shutdown hook always used.
|
||||
*
|
||||
* <p><strong>One statement could not move without reordering startup.</strong> {@code Fleetd.main}
|
||||
* registered its shutdown hook <em>before</em> building the Javalin {@code FleetApp} — the hook
|
||||
* itself never touched {@code app} (it still doesn't; see {@link FleetdRuntime#close()}), but the
|
||||
* hook needs a {@link FleetdRuntime} to close over, and the ticket asks that runtime to also own
|
||||
* {@code FleetApp}. Building the runtime before the app exists and mutating it afterward (via
|
||||
* {@link FleetdRuntime#attachApp}) preserves the exact original order — hook registered, then app
|
||||
* built, then HTTP started — without moving the app's construction earlier or the hook's
|
||||
* registration later. That is the one seam this ticket did not get to pin any other way.
|
||||
*
|
||||
* <p><strong>No inert variant.</strong> Deliberately, there is no overload of this method that
|
||||
* accepts a smaller/optional {@link ResourcePorts} or defaults one internally. A future edit that
|
||||
* wants to skip {@code FleetdAssembly} entirely and build its own graph is still possible — no
|
||||
* static analysis stops that — but it cannot do so by quietly swapping this call for an inert
|
||||
* substitute that still compiles, because none exists.
|
||||
*/
|
||||
final class FleetdAssembly {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(Fleetd.class);
|
||||
|
||||
/** CB-637: how often the lead coordination loop looks for peer messages — see {@code Fleetd}'s own constant. */
|
||||
private static final long LEAD_COORD_INTERVAL_MS = 3_000L;
|
||||
|
||||
private FleetdAssembly() {
|
||||
}
|
||||
|
||||
static FleetdRuntime assembleAndStart(AssemblyInputs inputs, ResourcePorts ports) {
|
||||
FleetConfig cfg = inputs.cfg();
|
||||
ConfigRef config = inputs.config();
|
||||
SubscriptionGuard guard = inputs.guard();
|
||||
|
||||
// fleetd #612 A-gaps (gap 2): moved in from Fleetd.main, immediately after cfg.validateAll()
|
||||
// there — the exact point main used to call these two, and still the first thing this
|
||||
// method does, before any socket or broker work below. See this class's javadoc and each
|
||||
// method's own for why they run here rather than in FleetConfig#validateAll() itself.
|
||||
Fleetd.reportRoleFallbackGaps(cfg);
|
||||
Fleetd.assertChartersNameOnlyRegisteredTools(cfg);
|
||||
|
||||
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
|
||||
? Path.of(cfg.herdrSocket())
|
||||
: UnixSocketHerdrClient.defaultSocketPath();
|
||||
|
||||
HerdrClient herdr = ports.connectHerdr(socket);
|
||||
HerdrClient memberHerdr = cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank()
|
||||
? ports.connectHerdr(Path.of(cfg.memberHerdrSocket()))
|
||||
: herdr;
|
||||
AtomicReference<Supplier<Map<String, String>>> leadsRef = new AtomicReference<>(Map::of);
|
||||
HerdrRouter router = new HerdrRouter(herdr, memberHerdr,
|
||||
target -> leadsRef.get().get().containsKey(target));
|
||||
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
|
||||
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
|
||||
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
|
||||
Map<String, FleetConfig.Profile> claudeProfiles = new LinkedHashMap<>();
|
||||
Map<String, FleetConfig.Profile> opencodeProfiles = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, w) -> {
|
||||
if (w.isOpenCode()) {
|
||||
opencodeProfiles.put(name, w);
|
||||
} else {
|
||||
claudeProfiles.put(name, w);
|
||||
}
|
||||
});
|
||||
List<HerdrPeerLauncher> adapters = new ArrayList<>();
|
||||
// fleetd #175: the daemon's real ExhaustionSink can only be built once `sessions` exists
|
||||
// (below), but `sessions` needs `workers`, which needs the adapters built right here — a
|
||||
// genuine cycle. Break it exactly like liveCountRef below: a forwarding sink built now,
|
||||
// pointed at the real one once it exists.
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
ExhaustionSink forwardingExhaustionSink = Fleetd.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.
|
||||
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
|
||||
adapters.add(Fleetd.claudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
|
||||
claudeProfiles, cfg, config));
|
||||
}
|
||||
if (!opencodeProfiles.isEmpty()) {
|
||||
adapters.add(Fleetd.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
|
||||
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
|
||||
BackendQuarantine quarantine = BackendQuarantine.withEscalation(ports.nanoClock(),
|
||||
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
|
||||
// fleetd #201 Unit 5: one outage-cool-off tracker for the whole daemon, shared between the
|
||||
// launcher (checked at spawn, like `quarantine` above) and the backend-error sink wired in
|
||||
// below (written on a classified backend error).
|
||||
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(ports.nanoClock());
|
||||
PeerLauncher workers = new CompositePeerLauncher(
|
||||
adapters,
|
||||
cfg.effectiveDefaultProfile(),
|
||||
config,
|
||||
profileName -> liveCountRef.get().apply(profileName),
|
||||
quarantine,
|
||||
outagePolicy);
|
||||
// fleetd #422 follow-up: say which of the three model-gate states the daemon booted into.
|
||||
log.info("model gate (fleetd #422): {}", Fleetd.modelGateCoverageLine(workers.modelGateState()));
|
||||
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
|
||||
// exists. Wait, then degrade rather than die: serving with /healthz reporting "degraded" is
|
||||
// strictly more useful than exiting.
|
||||
Fleetd.HerdrAwaitOutcome herdrOutcome = Fleetd.awaitHerdr(herdr, ports.nanoClock(), Fleetd::sleepHerdrPoll);
|
||||
boolean herdrUp = Fleetd.logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
|
||||
if (herdrUp) {
|
||||
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
|
||||
// with the previous process — reap those leaked orphans now, before we start serving.
|
||||
workers.reapOrphanWorkers();
|
||||
}
|
||||
|
||||
// CB-301: authoritative session registry + lifecycle FSM on top of ClaudeCodeLauncher.
|
||||
// CB-301-ext: worktree provisioning seam, optionally rooted at a configured directory.
|
||||
// CB-303 part 2: context cap is opt-in and disabled (0) when absent/null.
|
||||
int contextCap = 0;
|
||||
if (cfg.lifecycle() != null && cfg.lifecycle().contextCap() != null
|
||||
&& cfg.lifecycle().contextCap() > 0) {
|
||||
contextCap = cfg.lifecycle().contextCap();
|
||||
}
|
||||
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
|
||||
SessionManager sessions = new SessionManager(workers,
|
||||
new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup(), cfg.memberSkills()),
|
||||
ports.nanoClock(), contextCap, clearAfterTurn);
|
||||
liveCountRef.set(profileName -> Fleetd.liveSessionCount(sessions.roster(), profileName));
|
||||
|
||||
// Idle-sleep guard: hold an OS-level assertion against idle sleep while at least one
|
||||
// member is live. No-op (never constructed) off macOS or when idleSleepGuard.enabled is
|
||||
// explicitly false; the mechanism itself is additionally a no-op if 'caffeinate' cannot be
|
||||
// started, so this can never fail a spawn, a release, or startup.
|
||||
boolean idleSleepGuardEnabled = cfg.idleSleepGuard() == null || cfg.idleSleepGuard().isEnabled();
|
||||
final IdleSleepGuard idleSleepGuard;
|
||||
if (idleSleepGuardEnabled) {
|
||||
idleSleepGuard = new IdleSleepGuard(new CaffeinateSleepAssertionMechanism(), sessions::size);
|
||||
sessions.onAcquire(_ -> idleSleepGuard.recheck());
|
||||
sessions.onRelease(_ -> idleSleepGuard.recheck());
|
||||
} else {
|
||||
idleSleepGuard = null;
|
||||
}
|
||||
|
||||
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled. FIRST of the
|
||||
// recurring background loops to start — see this class's own javadoc for the full order.
|
||||
final SessionReaper reaper;
|
||||
if (cfg.lifecycle() != null
|
||||
&& cfg.lifecycle().idleTtlSeconds() != null
|
||||
&& cfg.lifecycle().idleTtlSeconds() > 0) {
|
||||
reaper = new SessionReaper(sessions, cfg.lifecycle().idleTtlSeconds());
|
||||
reaper.start();
|
||||
} else {
|
||||
reaper = null;
|
||||
}
|
||||
|
||||
// CB-530: every pane the config names as a lead, merged from `leaders:` and the legacy
|
||||
// singular pin. PrimaryRegistry below still tracks ONE terminal — it addresses the push
|
||||
// loop's nudges, which need a single destination — so it keeps the legacy pin.
|
||||
Map<String, String> leadTerminals = cfg.leaderTerminals();
|
||||
if (leadTerminals.size() > 1) {
|
||||
log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values());
|
||||
}
|
||||
// CB-531/CB-579: discover leads by the tab labels the operator writes, one scanner per
|
||||
// configured lead's own exact `tab:` label.
|
||||
final Supplier<Map<String, String>> leads;
|
||||
var leaders = cfg.fleet().leaders();
|
||||
if (!leaders.isEmpty()) {
|
||||
Map<String, String> tabToName = new LinkedHashMap<>();
|
||||
leaders.forEach((name, leader) -> {
|
||||
if (leader != null && leader.tab() != null && !leader.tab().isBlank()) {
|
||||
tabToName.put(leader.tab(), name);
|
||||
}
|
||||
});
|
||||
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
|
||||
// This must use the lead daemon: scanning member tabs would demote the lead to a worker.
|
||||
leads = new LeadTabScanner(herdr, tabToName, Set.of(),
|
||||
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), ports.nanoClock());
|
||||
log.info("lead scan: tabs {} host a lead (rescan every {}s, shared fleet space)",
|
||||
tabToName.keySet(), scanIntervalSeconds);
|
||||
} else {
|
||||
leads = () -> leadTerminals;
|
||||
}
|
||||
leadsRef.set(leads);
|
||||
|
||||
// CB-558: start any declared lead that is not already running. After the scanner is built,
|
||||
// and only when herdr answered — the launcher's whole safety property is that it can count
|
||||
// live leads first, and must never guess and risk a second orchestrator.
|
||||
if (herdrUp && !leaders.isEmpty()) {
|
||||
int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads();
|
||||
if (launched > 0) {
|
||||
log.info("lead auto-launch: {} lead(s) started", launched);
|
||||
}
|
||||
}
|
||||
|
||||
// CB-548: config-declared architect slots. Nothing here spawns a slot; the terminal → slot
|
||||
// binding is owned by the registry and empty at startup.
|
||||
MemberRegistry members = MemberRegistry.live(() -> config.get().fleet());
|
||||
sessions.setMemberLifecycle(members);
|
||||
if (!members.slots().isEmpty()) {
|
||||
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
|
||||
+ "spawn lifecycle binds a live terminal to it)",
|
||||
members.slots().size(), members.slots().keySet());
|
||||
}
|
||||
|
||||
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
|
||||
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
// CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured
|
||||
// usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer.
|
||||
LiveExhaustedPatterns liveExhaustedPatterns = Fleetd.liveExhaustedPatterns(config);
|
||||
ExhaustedPatternLookup exhaustedPatterns = Fleetd.exhaustedPatternLookup(sessions::roster, liveExhaustedPatterns);
|
||||
// The startup coverage line still reports the boot-time snapshot only.
|
||||
Set<String> exhaustedConfiguredAtStartup = cfg.profiles().entrySet().stream()
|
||||
.filter(e -> e.getValue().hasExhaustedPattern())
|
||||
.map(Map.Entry::getKey)
|
||||
.collect(Collectors.toCollection(LinkedHashSet::new));
|
||||
log.info("backend-exhausted classification (CB-578 stage A): {}",
|
||||
Fleetd.exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedConfiguredAtStartup));
|
||||
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
|
||||
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
|
||||
// rather than handing it back as a real answer.
|
||||
Map<String, Pattern> errorPatternsByProfile = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, profile) -> {
|
||||
if (profile.hasErrorPattern()) {
|
||||
errorPatternsByProfile.put(name, Pattern.compile(profile.errorPattern()));
|
||||
}
|
||||
});
|
||||
BackendErrorPatternLookup backendErrorPatterns = Fleetd.backendErrorPatternLookup(sessions::roster,
|
||||
errorPatternsByProfile);
|
||||
log.info("backend-error classification (fleetd #201 Unit 5): {}",
|
||||
Fleetd.errorPatternCoverageLine(cfg.profiles().keySet(), errorPatternsByProfile.keySet()));
|
||||
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
|
||||
// CREDENTIAL, so a profile sharing that credential is refused too, not just the one that
|
||||
// happened to report it.
|
||||
Map<String, String> quarantineReasonByCredential = new ConcurrentHashMap<>();
|
||||
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
|
||||
// now that `sessions` exists to resolve target -> session -> profile.
|
||||
ExhaustionSink exhaustionSink = Fleetd.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
|
||||
// holder set once `pushLoop` exists, read lazily from inside the lambda built here.
|
||||
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>();
|
||||
BackendErrorSink backendErrorSink = Fleetd.backendErrorSink(sessions, () -> config.get().profiles(),
|
||||
outagePolicy, pushLoopRef::get);
|
||||
AgentControl agents = router.memberAgents();
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns,
|
||||
exhaustionSink, backendErrorPatterns, backendErrorSink, ports.nanoClock(),
|
||||
Fleetd.worktreeBranchLookup(sessions::roster));
|
||||
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
||||
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
|
||||
MemberPresence presence = sessions.asPresence();
|
||||
TurnListener turnListener = Fleetd.turnListener(completion, sessions);
|
||||
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads);
|
||||
// fleetd #556: registration is wired directly to `completion`, not folded into the
|
||||
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
|
||||
// listener) throwing, regardless of call order.
|
||||
Injector injector = new Injector(router, turnListener, deliverable,
|
||||
presence::forget, Fleetd.turnRegistrar(completion));
|
||||
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||
poller.start(); // SECOND of the recurring background loops to start, after the reaper.
|
||||
|
||||
// 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 = Fleetd.selectReplyInbox(cfg.broker(), ports.environment(),
|
||||
ports.replyInboxOpener());
|
||||
// CB-637: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate
|
||||
// broker from the reply inbox by design. Absent a coordinator: block this is null and every
|
||||
// lead path below is simply not wired, exactly the behaviour before this ticket. It owns a
|
||||
// broker connection, so keep the reference for the ordered shutdown hook.
|
||||
final LeadChannelHandle leadMailbox = Fleetd.openLeadMailbox(cfg.coordinator(), ports.environment(),
|
||||
ports.leadMailboxOpener());
|
||||
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
|
||||
String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null;
|
||||
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal);
|
||||
// CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs.
|
||||
if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) {
|
||||
log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes "
|
||||
+ "from leaders:/leadScan:, and reply nudges follow the lead that delegated. "
|
||||
+ "It still works, and is still the fallback nudge destination when a restart "
|
||||
+ "has lost the delegation map. Its pushReminders/pushBackoffMs stay valid.");
|
||||
}
|
||||
// CB-307: active push-to-primary loop — nudge the primary when replies land without an
|
||||
// open fleet_send. Uses its own lightweight scheduled executor, separate from the injector.
|
||||
int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
|
||||
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
|
||||
var pushScheduler = ports.newScheduler("bridge-push-");
|
||||
// CB-502: the registry is built before the service and the push loop so send/reply outcomes
|
||||
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
|
||||
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
|
||||
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
|
||||
pushScheduler, maxReminders, backoffMs, metrics);
|
||||
// fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda
|
||||
// above at the real push loop, now that it exists.
|
||||
pushLoopRef.set(pushLoop);
|
||||
// CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so
|
||||
// an upgraded daemon cannot silently start spending subscription on nudging an idle lead.
|
||||
// THIRD of the recurring background loops to start (optional).
|
||||
final LeadHeartbeatLoop heartbeat;
|
||||
var heartbeatScheduler = ports.newScheduler("bridge-heartbeat-");
|
||||
if (cfg.leadHeartbeat() != null) {
|
||||
var hb = cfg.leadHeartbeat();
|
||||
var leadContextGauge = new LeadContextGauge();
|
||||
// fleetd #621: the context-high notice's own wording must track this same effective
|
||||
// value — LeadRollover.confirm(...) already gates the roll on it (LeadRollover.java:480),
|
||||
// and absent `leadRollover:` entirely the roll is unusable regardless (NOT_CONFIGURED),
|
||||
// so `true` (the FleetConfig.LeadRollover default) is the safe, byte-identical fallback.
|
||||
// Carried in from Fleetd.main when #612 Unit A merged main: #622 added this line to the
|
||||
// block Unit A had already moved here, so the merge would otherwise have silently
|
||||
// dropped it — with a fully green suite, because nothing pins it (see the follow-up issue).
|
||||
boolean requireOperatorConfirm = cfg.leadRollover() == null || cfg.leadRollover().requireOperatorConfirm();
|
||||
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
|
||||
pushLoop, heartbeatScheduler, ports.nanoClock(),
|
||||
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
|
||||
metrics,
|
||||
Fleetd.leadContextSource(leadContextGauge, router.leadAgents(), leads,
|
||||
Fleetd.leadConfigDirLookup(() -> config.get().profiles(), leaders)),
|
||||
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm);
|
||||
heartbeat.start();
|
||||
} else {
|
||||
heartbeat = null;
|
||||
heartbeatScheduler.shutdownNow();
|
||||
}
|
||||
// fleetd #480: lead rollover. Opt-in; absent `leadRollover:` this is never constructed.
|
||||
LeadRollover leadRollover = Fleetd.leadRollover(cfg, router.leadAgents(), config, leads);
|
||||
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
|
||||
pushLoop, metrics);
|
||||
|
||||
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
|
||||
// FOURTH of the recurring background loops to start (optional).
|
||||
final FleetHealthMonitor healthMonitor;
|
||||
var healthScheduler = ports.newScheduler("bridge-health-");
|
||||
if (cfg.health() != null && cfg.health().isEnabled()) {
|
||||
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it.
|
||||
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
|
||||
ports.nanoClock(), ports.wallClockNanos(),
|
||||
cfg.health().intervalOrDefault(),
|
||||
cfg.health().workingSuspectAfterOrDefault(), Fleetd.healthFailTarget(messages));
|
||||
String coverage = FleetHealthMonitor.coverage(true,
|
||||
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
||||
if ("detection-only".equals(coverage)) {
|
||||
log.warn("fleet health: {} (no notification sink configured)", coverage);
|
||||
} else {
|
||||
log.info("fleet health: {}", coverage);
|
||||
}
|
||||
healthMonitor.start();
|
||||
} else {
|
||||
healthMonitor = null;
|
||||
healthScheduler.shutdownNow();
|
||||
}
|
||||
|
||||
// CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
|
||||
// release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
|
||||
sessions.onAcquire(replyInbox::own);
|
||||
// CB-516: releasing a worker must fail whatever send was waiting on it.
|
||||
sessions.onRelease(Fleetd.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.
|
||||
ConnectionIdentity identity = new ConnectionIdentity(
|
||||
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||
|
||||
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
|
||||
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
|
||||
final CallerResolver callers;
|
||||
if (cfg.auth().tokenMode()) {
|
||||
String token = ports.environment().get(cfg.auth().tokenEnv());
|
||||
if (token == null || token.isBlank()) {
|
||||
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
|
||||
+ " is unset or empty — export it before starting fleetd");
|
||||
}
|
||||
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members);
|
||||
log.info("auth: token mode (bearer required for non-worker callers, env {})",
|
||||
cfg.auth().tokenEnv());
|
||||
} else {
|
||||
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members);
|
||||
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
|
||||
}
|
||||
|
||||
FleetMcp.QuarantineSource quarantineSource = Fleetd.quarantineSource(config, quarantine,
|
||||
liveExhaustedPatterns, quarantineReasonByCredential);
|
||||
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy);
|
||||
|
||||
FleetMcp.LoopHealthSource loopHealth = Fleetd.loopHealthSource(poller, reaper);
|
||||
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
|
||||
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
|
||||
Fleetd.capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
|
||||
Fleetd.healthCoverageSource(config),
|
||||
loopHealth,
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
new FleetMcp.LeadSeatSource(Fleetd.leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
|
||||
Fleetd.leadConfigDirSource(() -> config.get().profiles(), leaders),
|
||||
cfg.coordinator() == null ? List.of() : cfg.coordinator().peers(),
|
||||
leadRollover);
|
||||
|
||||
// CB-637: the receive half. Only constructed when a lead mailbox actually opened.
|
||||
final LeadCoordLoop leadCoordLoop;
|
||||
final ScheduledExecutorService leadCoordSchedulerRef;
|
||||
if (leadMailbox != null) {
|
||||
var leadCoordScheduler = ports.newScheduler("bridge-leadcoord-");
|
||||
leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler,
|
||||
LEAD_COORD_INTERVAL_MS);
|
||||
leadCoordLoop.start(); // FIFTH of the recurring background loops to start (optional).
|
||||
leadCoordSchedulerRef = leadCoordScheduler;
|
||||
} else {
|
||||
leadCoordLoop = null;
|
||||
leadCoordSchedulerRef = null;
|
||||
}
|
||||
|
||||
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed.
|
||||
// SIXTH of the recurring background loops to start (optional).
|
||||
final ConfigWatcher configWatcher;
|
||||
if (cfg.configReload() != null && cfg.configReload().isEnabled()) {
|
||||
configWatcher = new ConfigWatcher(config, cfg.configReload().intervalSeconds());
|
||||
configWatcher.start();
|
||||
} else {
|
||||
configWatcher = null;
|
||||
}
|
||||
|
||||
// CB-303 part 3: single ordered shutdown hook — see FleetdRuntime#close() for the statements
|
||||
// this used to be. Registered here, at the exact point `main` used to register it: after
|
||||
// configWatcher, before the Javalin app exists (see this class's own javadoc for why).
|
||||
FleetdRuntime runtime = new FleetdRuntime(cfg, sessions, router, poller, messages, pushLoop, heartbeat,
|
||||
leadCoordLoop, leadCoordSchedulerRef, healthMonitor, configWatcher, mcp, reaper, idleSleepGuard,
|
||||
replyInbox, leadMailbox, completion, injector);
|
||||
ports.addShutdownHook(runtime::close);
|
||||
|
||||
// CB-185: give FleetApp both daemons — /healthz must require both to answer and
|
||||
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
|
||||
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
|
||||
callers, metrics, deliverable,
|
||||
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
|
||||
quarantineSource, outageSource, loopHealth).build();
|
||||
runtime.attachApp(app);
|
||||
ports.startHttp(app, cfg.bind().host(), cfg.bind().port()); // HTTP starts LAST, always.
|
||||
log.info("fleetd listening on {}:{}, herdr socket {}",
|
||||
cfg.bind().host(), cfg.bind().port(), socket);
|
||||
return runtime;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,166 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigWatcher;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.health.FleetHealthMonitor;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.StatusPoller;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.msg.LeadChannelHandle;
|
||||
import dev.ltms.fleet.msg.LeadCoordLoop;
|
||||
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.msg.ReplyPushLoop;
|
||||
import dev.ltms.fleet.power.IdleSleepGuard;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.SessionReaper;
|
||||
import io.javalin.Javalin;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
|
||||
/**
|
||||
* fleetd #612 Unit A: the assembled daemon. {@link FleetdAssembly#assembleAndStart} builds exactly
|
||||
* one of these and hands it to {@code ResourcePorts.addShutdownHook}; {@link #close()} is the
|
||||
* single ordered shutdown, moved verbatim out of {@code Fleetd.main}'s old shutdown-hook
|
||||
* {@code Thread} body — same statements, same order, see that method's javadoc.
|
||||
*
|
||||
* <p><strong>Owns the real production objects, never a copy.</strong> Every package-private
|
||||
* accessor below returns the identical instance the running daemon is using. That is the entire
|
||||
* point of this class existing (see fleetd #612's problem statement): a test that inspected a
|
||||
* snapshot built alongside the real objects could pass while production silently received
|
||||
* something else — the exact shape of the #602/#606 defect this ticket exists to stop from
|
||||
* recurring one call site at a time. Nothing here is rebuilt or copied for a test's benefit.
|
||||
*/
|
||||
final class FleetdRuntime implements AutoCloseable {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(FleetdRuntime.class);
|
||||
|
||||
private final FleetConfig cfg;
|
||||
private final SessionManager sessions;
|
||||
private final HerdrRouter router;
|
||||
private final StatusPoller poller;
|
||||
private final MessageService messages;
|
||||
private final ReplyPushLoop pushLoop;
|
||||
private final LeadHeartbeatLoop heartbeat; // nullable — leadHeartbeat: opt-in
|
||||
private final LeadCoordLoop leadCoordLoop; // nullable — coordinator: opt-in
|
||||
private final ScheduledExecutorService leadCoordScheduler; // nullable, paired with leadCoordLoop
|
||||
private final FleetHealthMonitor healthMonitor; // nullable — health.enabled opt-in
|
||||
private final ConfigWatcher configWatcher; // nullable — configReload.enabled opt-in
|
||||
private final FleetMcp mcp;
|
||||
private final SessionReaper reaper; // nullable — lifecycle.idleTtlSeconds opt-in
|
||||
private final IdleSleepGuard idleSleepGuard; // nullable — idleSleepGuard.enabled: false
|
||||
private final ReplyInbox replyInbox;
|
||||
private final LeadChannelHandle leadMailbox; // nullable — coordinator: opt-in
|
||||
private final CompletionResolver completion;
|
||||
private final Injector injector;
|
||||
/**
|
||||
* Not final: {@code Fleetd.main}'s shutdown hook was registered <em>before</em> the Javalin
|
||||
* {@code FleetApp} was built and started — see {@link FleetdAssembly#assembleAndStart}'s javadoc
|
||||
* for why that order could not be preserved AND have this constructor take {@code app}. {@link
|
||||
* #attachApp} is called immediately after the real app is built, still before HTTP starts
|
||||
* listening, so this is set long before any test or caller could observe it unset.
|
||||
*/
|
||||
private Javalin app;
|
||||
|
||||
FleetdRuntime(FleetConfig cfg, SessionManager sessions, HerdrRouter router, StatusPoller poller,
|
||||
MessageService messages, ReplyPushLoop pushLoop, LeadHeartbeatLoop heartbeat,
|
||||
LeadCoordLoop leadCoordLoop, ScheduledExecutorService leadCoordScheduler,
|
||||
FleetHealthMonitor healthMonitor, ConfigWatcher configWatcher, FleetMcp mcp,
|
||||
SessionReaper reaper, IdleSleepGuard idleSleepGuard, ReplyInbox replyInbox,
|
||||
LeadChannelHandle leadMailbox, CompletionResolver completion, Injector injector) {
|
||||
this.cfg = cfg;
|
||||
this.sessions = sessions;
|
||||
this.router = router;
|
||||
this.poller = poller;
|
||||
this.messages = messages;
|
||||
this.pushLoop = pushLoop;
|
||||
this.heartbeat = heartbeat;
|
||||
this.leadCoordLoop = leadCoordLoop;
|
||||
this.leadCoordScheduler = leadCoordScheduler;
|
||||
this.healthMonitor = healthMonitor;
|
||||
this.configWatcher = configWatcher;
|
||||
this.mcp = mcp;
|
||||
this.reaper = reaper;
|
||||
this.idleSleepGuard = idleSleepGuard;
|
||||
this.replyInbox = replyInbox;
|
||||
this.leadMailbox = leadMailbox;
|
||||
this.completion = completion;
|
||||
this.injector = injector;
|
||||
}
|
||||
|
||||
/** See the {@link #app} field doc for why this is a late-bound setter rather than a constructor arg. */
|
||||
void attachApp(Javalin app) {
|
||||
this.app = app;
|
||||
}
|
||||
|
||||
// --- package-private accessors: the SAME instances this runtime owns, never a copy ---------
|
||||
|
||||
SessionManager sessions() { return sessions; }
|
||||
HerdrRouter router() { return router; }
|
||||
StatusPoller poller() { return poller; }
|
||||
MessageService messages() { return messages; }
|
||||
ReplyPushLoop pushLoop() { return pushLoop; }
|
||||
LeadHeartbeatLoop heartbeat() { return heartbeat; }
|
||||
LeadCoordLoop leadCoordLoop() { return leadCoordLoop; }
|
||||
FleetHealthMonitor healthMonitor() { return healthMonitor; }
|
||||
ConfigWatcher configWatcher() { return configWatcher; }
|
||||
FleetMcp mcp() { return mcp; }
|
||||
SessionReaper reaper() { return reaper; }
|
||||
IdleSleepGuard idleSleepGuard() { return idleSleepGuard; }
|
||||
ReplyInbox replyInbox() { return replyInbox; }
|
||||
LeadChannelHandle leadMailbox() { return leadMailbox; }
|
||||
CompletionResolver completion() { return completion; }
|
||||
Injector injector() { return injector; }
|
||||
Javalin app() { return app; }
|
||||
|
||||
/**
|
||||
* CB-303 part 3: the single ordered shutdown. Moved verbatim out of {@code Fleetd.main}'s
|
||||
* shutdown-hook {@code Thread} body (fleetd #612 Unit A) — drain sessions first while herdr is
|
||||
* still open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close
|
||||
* herdr last, exactly as before. {@code Fleetd.main} never calls this directly; it hands the
|
||||
* reference to {@code ResourcePorts.addShutdownHook} the moment this runtime exists, the same
|
||||
* point it used to register the hook {@code Thread} itself.
|
||||
*/
|
||||
@Override
|
||||
public void close() {
|
||||
sessions.close(cfg.lifecycle() != null ? cfg.lifecycle().drainTimeoutSeconds() : null);
|
||||
poller.stop();
|
||||
messages.close();
|
||||
pushLoop.close();
|
||||
if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler
|
||||
if (leadCoordLoop != null) leadCoordLoop.close(); // CB-637: stop delivering peer-lead messages
|
||||
if (leadCoordScheduler != null) leadCoordScheduler.shutdownNow();
|
||||
if (healthMonitor != null) healthMonitor.stop();
|
||||
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
|
||||
mcp.close();
|
||||
if (reaper != null) reaper.stop();
|
||||
// Idle-sleep guard: release unconditionally, even though sessions.close() above already
|
||||
// drained every session (and each release already drove the live count to 0, which
|
||||
// releases the guard's assertion on its own) — this is the backstop for a drain that was
|
||||
// itself interrupted or threw, so no caffeinate child ever outlives the daemon.
|
||||
if (idleSleepGuard != null) idleSleepGuard.close();
|
||||
// Release the broker connection last among message resources (no-op for the in-memory inbox).
|
||||
if (replyInbox instanceof AutoCloseable closeable) {
|
||||
try {
|
||||
closeable.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("reply inbox close: {}", e.toString());
|
||||
}
|
||||
}
|
||||
// CB-637: the coordination connection goes with it — after the loop that reads it has
|
||||
// stopped, so no tick can be mid-ack against a closed channel.
|
||||
if (leadMailbox != null) {
|
||||
try {
|
||||
leadMailbox.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("lead mailbox close: {}", e.toString());
|
||||
}
|
||||
}
|
||||
router.close();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import io.javalin.Javalin;
|
||||
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
/**
|
||||
* fleetd #612 Unit A: every boot-time side effect {@link FleetdAssembly#assembleAndStart} performs
|
||||
* that a real daemon must do for real, and a test must not — read the process environment, connect
|
||||
* a herdr client, open a broker (the reply inbox, the lead mailbox), read a clock, start a
|
||||
* background scheduler, register the JVM shutdown hook, and bind the HTTP server.
|
||||
*
|
||||
* <p>{@link #system()} is the one production implementation ({@code SystemResourcePorts}), wired
|
||||
* verbatim from what {@code Fleetd.main} used to call directly at each of these call sites. A test
|
||||
* builds its own implementation instead of receiving an inert default from this interface —
|
||||
* deliberately, there is no {@code ResourcePorts.none()}. fleetd #612's whole problem is a call
|
||||
* site quietly swapped for an inert variant that still compiles; adding one here, even for tests,
|
||||
* would hand a future edit to {@code FleetdAssembly} the exact compiling substitute this ticket
|
||||
* exists to rule out. A test that wants an inert resource writes its own fake and owns that
|
||||
* decision explicitly.
|
||||
*/
|
||||
public interface ResourcePorts {
|
||||
|
||||
/** The process environment. Production: {@link System#getenv()}. */
|
||||
Map<String, String> environment();
|
||||
|
||||
/**
|
||||
* Connect a herdr client bound to {@code socketPath}. Production returns a real
|
||||
* {@code UnixSocketHerdrClient} — connection-per-call, so this itself never touches the socket.
|
||||
*/
|
||||
HerdrClient connectHerdr(Path socketPath);
|
||||
|
||||
/** The reply-inbox AMQP opener (CB-307). Production: {@link Fleetd#replyInboxOpener()}. */
|
||||
Fleetd.AmqpOpener replyInboxOpener();
|
||||
|
||||
/** The lead-mailbox AMQP opener (CB-637). Production: {@link Fleetd#leadMailboxOpener()}. */
|
||||
Fleetd.LeadMailboxOpener leadMailboxOpener();
|
||||
|
||||
/** A monotonic elapsed-time clock. Production: {@link System#nanoTime()}. */
|
||||
LongSupplier nanoClock();
|
||||
|
||||
/**
|
||||
* A wall-clock reading, in nanoseconds. Production: {@code System.currentTimeMillis()}
|
||||
* converted to nanoseconds. Kept separate from {@link #nanoClock()} because {@link
|
||||
* dev.ltms.fleet.health.FleetHealthMonitor} needs both — one monotonic clock for elapsed-time
|
||||
* decisions, one wall clock to detect and correct for a macOS sleep freezing the monotonic one.
|
||||
*/
|
||||
LongSupplier wallClockNanos();
|
||||
|
||||
/** A dedicated single-thread scheduler; production names its (virtual) thread {@code purpose}. */
|
||||
ScheduledExecutorService newScheduler(String purpose);
|
||||
|
||||
/** Register a JVM shutdown hook that runs {@code hook} on JVM exit. */
|
||||
void addShutdownHook(Runnable hook);
|
||||
|
||||
/** Bind and start the HTTP server. */
|
||||
void startHttp(Javalin app, String host, int port);
|
||||
|
||||
/** The real production ports: a live herdr socket, a live broker, real threads, a real bind. */
|
||||
static ResourcePorts system() {
|
||||
return new SystemResourcePorts();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.UnixSocketHerdrClient;
|
||||
import io.javalin.Javalin;
|
||||
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
/**
|
||||
* fleetd #612 Unit A: the one production {@link ResourcePorts} — every method here is the exact
|
||||
* call {@code Fleetd.main} used to make directly at each of these sites before this ticket.
|
||||
* Package-private: obtained only through {@link ResourcePorts#system()}.
|
||||
*/
|
||||
final class SystemResourcePorts implements ResourcePorts {
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return System.getenv();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return UnixSocketHerdrClient.connect(socketPath, new ObjectMapper());
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return Fleetd.replyInboxOpener();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return Fleetd.leadMailboxOpener();
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis());
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor(r -> Thread.ofVirtual().name(purpose).unstarted(r));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(hook));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
app.start(host, port);
|
||||
}
|
||||
}
|
||||
@@ -31,6 +31,20 @@ public final class ConnectionIdentity {
|
||||
this.cwds = cwds;
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@link PaneLocator} this identity resolves callers against — fleetd #612 CB-185: lets a
|
||||
* test drive the exact {@link PaneLocator} a real assembly wired up (e.g. {@code
|
||||
* FleetdAssembly}'s {@code new ConnectionIdentity(new PaneLocator(herdr, memberHerdr), ...)})
|
||||
* directly with a chosen pid, bypassing the OS-dependent {@link PeerPidLookup} that {@link
|
||||
* #resolve} otherwise goes through. A full HTTP round trip cannot exercise this: {@code
|
||||
* LsofPeerPidLookup} excludes its own pid, and an in-process test client and server share one
|
||||
* JVM pid, so {@code pidForLocalPort} always returns {@code -1} and {@link PaneLocator} never
|
||||
* gets called at all.
|
||||
*/
|
||||
public PaneLocator panes() {
|
||||
return panes;
|
||||
}
|
||||
|
||||
/**
|
||||
* The caller resolved from the connection: its worker {@code terminal} (or {@code null} for the
|
||||
* primary / an off-host client), its {@code pid} (or {@code -1} if not resolvable), and whether
|
||||
|
||||
@@ -107,6 +107,13 @@ public final class FleetMcp {
|
||||
* untested identity heuristic.
|
||||
*/
|
||||
private final boolean authorizationEnforced;
|
||||
/**
|
||||
* fleetd #612 CB-185: kept as a field (rather than only captured by the {@code
|
||||
* contextExtractor} closure built in the constructor) so a test can reach the exact {@link
|
||||
* ConnectionIdentity} — and, through {@link ConnectionIdentity#panes()}, the exact {@link
|
||||
* dev.ltms.fleet.herdr.PaneLocator} — that a real assembly wired up. See {@link #identity()}.
|
||||
*/
|
||||
private final ConnectionIdentity identity;
|
||||
private final Metrics metrics; // CB-502: null → auth failures not counted
|
||||
private final CapacitySource capacity;
|
||||
private final HealthCoverageSource healthCoverage;
|
||||
@@ -396,6 +403,7 @@ public final class FleetMcp {
|
||||
Objects.requireNonNull(callers, "callers");
|
||||
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
|
||||
== AuthorizationMode.ENFORCED;
|
||||
this.identity = identity;
|
||||
this.leadChannel = leadChannel;
|
||||
this.peers = peers == null ? List.of() : List.copyOf(peers);
|
||||
this.capacity = capacity;
|
||||
@@ -755,6 +763,15 @@ public final class FleetMcp {
|
||||
return transport;
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@link ConnectionIdentity} this server resolves every caller against — fleetd #612
|
||||
* CB-185: lets a test reach the exact {@link dev.ltms.fleet.herdr.PaneLocator} a real assembly
|
||||
* wired up (via {@link ConnectionIdentity#panes()}), rather than a copy built for the test.
|
||||
*/
|
||||
public ConnectionIdentity identity() {
|
||||
return identity;
|
||||
}
|
||||
|
||||
/** Mark a connected spawned member available for the injector readiness gate. */
|
||||
static void markSpawnedMemberPresent(Principal caller, MemberPresence presence) {
|
||||
if (caller.isSpawnedMember()) {
|
||||
@@ -780,6 +797,46 @@ public final class FleetMcp {
|
||||
return server.listTools();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — same reason as {@link #registeredTools()}: a test that must drive the REAL
|
||||
* {@link QuarantineSource} (and the real {@link BackendQuarantine} it wraps) this daemon was
|
||||
* assembled with, rather than scraping {@code FleetdAssembly.java}'s source text for the
|
||||
* constructor call that built it. Unlike {@link #registeredTools()}'s callers, that test cannot
|
||||
* live in this package: it also builds the {@code ResourcePorts} that drives
|
||||
* {@code FleetdAssembly.assembleAndStart}, and {@code ResourcePorts}' methods return
|
||||
* {@code Fleetd}-nested types that are only visible from package {@code dev.ltms.fleet} — so
|
||||
* this accessor is {@code public}, not package-private, to stay reachable from there. {@code
|
||||
* FleetdBackendQuarantineAssemblyTest} quarantines a credential twice through this exact
|
||||
* instance and checks the second cooldown is longer than the first — the one behavioural
|
||||
* difference {@link BackendQuarantine#withEscalation} and the flat two-argument constructor
|
||||
* actually produce.
|
||||
*/
|
||||
public QuarantineSource quarantineSource() {
|
||||
return quarantine;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — as {@link #quarantineSource()}, {@code public} for the same cross-package
|
||||
* reason, for the real {@link LeadSeatSource} this daemon was assembled with. {@code
|
||||
* FleetdLeadSeatAssemblyTest} calls {@code seatsFor} on this exact instance and checks it
|
||||
* reports a live lead's seat, which {@link LeadSeatSource#none()} can never do (it is a
|
||||
* constant-zero function regardless of input).
|
||||
*/
|
||||
public LeadSeatSource leadSeatSource() {
|
||||
return leadSeats;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — as {@link #quarantineSource()}, {@code public} for the same cross-package
|
||||
* reason, for the real {@link LeadRollover} (or {@code null}) this daemon was assembled with.
|
||||
* {@code FleetdLeadRolloverAssemblyTest} drives {@code open}/{@code confirm} on this exact
|
||||
* instance and waits for the real continuation to send {@code /clear} and {@code bootstrapText}
|
||||
* through the real {@code router.leadAgents()}.
|
||||
*/
|
||||
public LeadRollover leadRollover() {
|
||||
return leadRollover;
|
||||
}
|
||||
|
||||
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
/**
|
||||
* fleetd #612 A-gaps (gap 1): a {@link LeadChannel} that its owner can also close.
|
||||
*
|
||||
* <p>{@link LeadChannel}'s own javadoc says plainly that {@code close()} is deliberately left out
|
||||
* of that interface — draining is a caller convenience nobody uses, and closing is the
|
||||
* <em>owner's</em> job. This interface is that owner's own, wider view: whoever opens the
|
||||
* coordination mailbox (the assembly that builds the daemon) also needs to close it from the
|
||||
* shutdown path, and a test standing in for a real broker connection needs a fake it can mark
|
||||
* closed, without ever holding a live connection. Every ordinary consumer ({@code FleetMcp},
|
||||
* {@link LeadCoordLoop}) keeps taking the narrower {@link LeadChannel} exactly as before — only
|
||||
* the owner speaks this wider one.
|
||||
*
|
||||
* <p>{@link LeadMailbox} is still the only production implementation. This only generalises the
|
||||
* TYPE its owner holds it as (previously the concrete class), so a test can substitute a fake
|
||||
* closeable channel instead of a real AMQP connection.
|
||||
*/
|
||||
public interface LeadChannelHandle extends LeadChannel, AutoCloseable {
|
||||
|
||||
/** Release the underlying connection. Declared with no checked exception, unlike the plain {@link AutoCloseable#close()}. */
|
||||
@Override
|
||||
void close();
|
||||
}
|
||||
@@ -76,6 +76,7 @@ public final class LeadHeartbeatLoop {
|
||||
private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests
|
||||
private final LeadContextSource contextSource; // fleetd #609
|
||||
private final boolean contextHighNudge; // fleetd #609: opt-in, like the loop itself
|
||||
private final boolean requireOperatorConfirm; // fleetd #621: mirrors leadRollover.requireOperatorConfirm
|
||||
|
||||
/** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */
|
||||
private long idleSinceNanos = NOT_IDLE;
|
||||
@@ -106,12 +107,32 @@ public final class LeadHeartbeatLoop {
|
||||
* fleetd #609: as above, plus the lead's own context source and whether a HIGH reading should
|
||||
* append a hand-over notice to the loop's nudge. Pass {@link LeadContextSource#none()} and
|
||||
* {@code false} to keep the pre-#609 behaviour exactly (both existing public constructors do).
|
||||
*
|
||||
* <p>fleetd #621: delegates to the full constructor with {@code requireOperatorConfirm=true} —
|
||||
* the pre-#621 wording ("ask the operator ... only the operator can approve the roll") assumed
|
||||
* the config default, so every caller of this overload keeps that text byte-identical.
|
||||
*/
|
||||
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||
LeadContextSource contextSource, boolean contextHighNudge) {
|
||||
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
|
||||
idleAfterNanos, backoffMs, quietNudgeCap, metrics, contextSource, contextHighNudge, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #621: as above, plus the daemon's effective {@code leadRollover.requireOperatorConfirm}
|
||||
* value — threaded into {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean, boolean)}
|
||||
* so the notice's wording tracks the config the daemon actually enforces (see {@code
|
||||
* LeadRollover.confirm}) instead of always asserting the operator gate is on.
|
||||
*/
|
||||
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||
LeadContextSource contextSource, boolean contextHighNudge,
|
||||
boolean requireOperatorConfirm) {
|
||||
this.primaryRegistry = primaryRegistry;
|
||||
this.agents = agents;
|
||||
this.inbox = inbox;
|
||||
@@ -125,6 +146,7 @@ public final class LeadHeartbeatLoop {
|
||||
this.metrics = metrics;
|
||||
this.contextSource = contextSource;
|
||||
this.contextHighNudge = contextHighNudge;
|
||||
this.requireOperatorConfirm = requireOperatorConfirm;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -344,7 +366,7 @@ public final class LeadHeartbeatLoop {
|
||||
// d.contextNotified() is the value to persist once delivery is confirmed, not the value the text
|
||||
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
|
||||
// notice at all, defeating the very check this fixes.
|
||||
String notice = contextNotice(contextHighNudge, reading, contextNotified);
|
||||
String notice = contextNotice(contextHighNudge, reading, contextNotified, requireOperatorConfirm);
|
||||
var lead = primaryRegistry.primaryTerminal();
|
||||
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
|
||||
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
|
||||
@@ -381,8 +403,9 @@ public final class LeadHeartbeatLoop {
|
||||
/**
|
||||
* fleetd #609: the text appended to a nudge when the lead's own context is full — {@code ""}
|
||||
* whenever the notice does not apply, so callers can unconditionally append this without an extra
|
||||
* branch. Wording stays plain (CEFR B1) and honest that only the operator approves a roll — this
|
||||
* loop only ever prints text, it never calls {@code fleet_handover} itself.
|
||||
* branch. Wording stays plain (CEFR B1) and honest about who actually gates the roll — see the
|
||||
* {@code requireOperatorConfirm} overload (fleetd #621) for which check that is. This loop only
|
||||
* ever prints text, it never calls {@code fleet_handover} itself.
|
||||
*
|
||||
* @param enabled the {@code leadHeartbeat.contextHighNudge} config flag
|
||||
* @param reading the lead's current {@link LeadContextGauge} reading
|
||||
@@ -401,9 +424,33 @@ public final class LeadHeartbeatLoop {
|
||||
* closing sentence ("You will not be told again until your context reads ok.") false. {@link
|
||||
* #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}.
|
||||
*
|
||||
* <p>fleetd #621: delegates with {@code requireOperatorConfirm=true} — the pre-#621 default and the
|
||||
* value every existing caller of this overload (including every test written before #621) already
|
||||
* assumed, so the text this overload returns stays byte-identical.
|
||||
*
|
||||
* @param alreadyNotified whether the lead has already been told about the current HIGH stretch
|
||||
*/
|
||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) {
|
||||
return contextNotice(enabled, reading, alreadyNotified, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #621: as {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean)}, but the closing
|
||||
* instructions also track the daemon's effective {@code leadRollover.requireOperatorConfirm} value,
|
||||
* instead of always asserting that only the operator can approve the roll.
|
||||
*
|
||||
* <p>{@code LeadRollover.confirm(...)} already honours this flag: when it is {@code false}, the daemon
|
||||
* itself gates the roll on the three handover-file checks alone (exists, modified after the {@code
|
||||
* open()} request, and no older than {@code maxDocAgeSeconds}) and never consults {@code
|
||||
* operatorConfirmed}. Before this parameter existed, this notice told the lead to ask the operator
|
||||
* regardless — so a lead that followed its own instructions asked anyway, and setting the config knob
|
||||
* to {@code false} stopped the daemon refusing the roll without stopping the operator being
|
||||
* interrupted. This parameter is how the text is kept honest about which gate is actually live.
|
||||
*
|
||||
* @param requireOperatorConfirm the effective {@code leadRollover.requireOperatorConfirm} value
|
||||
*/
|
||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified,
|
||||
boolean requireOperatorConfirm) {
|
||||
if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) {
|
||||
return "";
|
||||
}
|
||||
@@ -420,10 +467,18 @@ public final class LeadHeartbeatLoop {
|
||||
sb.append(" (").append(reading.compactions()).append(' ').append(compactionWord)
|
||||
.append(" so far).");
|
||||
}
|
||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
|
||||
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
|
||||
+ "again until your context reads ok.");
|
||||
if (requireOperatorConfirm) {
|
||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
|
||||
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
|
||||
+ "again until your context reads ok.");
|
||||
} else {
|
||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||
+ "write the file it names, then call fleet_handover(action=\"confirm\", token). Decide for "
|
||||
+ "yourself when to confirm: the roll goes through if the handover file exists, was "
|
||||
+ "changed after you opened it, and is not older than maxDocAgeSeconds. You will not be "
|
||||
+ "told again until your context reads ok.");
|
||||
}
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
|
||||
@@ -63,7 +63,7 @@ import java.util.concurrent.TimeoutException;
|
||||
* which messages reached a lead. Any publish still awaiting its confirm is failed rather than left to idle out
|
||||
* the confirm timeout against a sequence number that means nothing on the new channel.
|
||||
*/
|
||||
public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
public final class LeadMailbox implements LeadChannelHandle {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(LeadMailbox.class);
|
||||
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
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.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
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.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 2, unit B2 (CB-185, identity half). Replaces the deleted
|
||||
* {@code FleetdConnectionIdentityConstructionTest}, which pinned this claim by reading {@code
|
||||
* Fleetd.java}'s source text for {@code "new PaneLocator(herdr, memberHerdr)"}. That claim moved
|
||||
* to {@code FleetdAssembly.java} (fleetd #612 Unit A) and is pinned here instead, by driving the
|
||||
* real {@link ConnectionIdentity} — via {@code runtime.mcp().identity()}, not a copy — that {@link
|
||||
* FleetdAssembly#assembleAndStart} built.
|
||||
*
|
||||
* <p><strong>What this guards against</strong> (from the deleted test's own javadoc): pinning
|
||||
* {@code PaneLocator} to {@code memberHerdr} alone leaves every LEAD's own MCP connection
|
||||
* unresolvable ({@code callerTerminal == null}) the moment {@code memberHerdrSocket} names a
|
||||
* second daemon, which breaks {@code fleet_reply}/{@code fleet_ask}/{@code fleet_whoami} for a
|
||||
* lead. {@code PaneLocatorTest} already proves {@link PaneLocator} itself can search two clients
|
||||
* given two — the gap this pins is that the assembly actually passes it two, and in the right
|
||||
* order (lead first).
|
||||
*
|
||||
* <p><strong>Why this cannot be driven through a real MCP/HTTP round trip.</strong> The natural
|
||||
* way to observe {@code ConnectionIdentity} would be a real {@code fleet_whoami} call over the
|
||||
* built {@code FleetMcp}, the way {@code FleetMcpContextExtractorTest} drives its own
|
||||
* hand-built one. That does not work for the REAL assembly, because {@code FleetdAssembly} wires
|
||||
* {@code ConnectionIdentity} with a hardcoded {@code new LsofPeerPidLookup()} (see {@code
|
||||
* FleetdAssembly.java:444}), and {@code LsofPeerPidLookup} explicitly excludes its own PID — see
|
||||
* its javadoc: "we exclude our own PID and take the other end". In a JUnit test the HTTP client
|
||||
* and the daemon under test run in the very same JVM, so the "client" and "server" ends of the
|
||||
* loopback connection ARE the same PID, and {@code pidForLocalPort} always returns {@code -1}
|
||||
* before {@link PaneLocator} is ever reached — proving nothing about which daemon(s) got searched.
|
||||
* This test instead reaches the real {@link PaneLocator} the assembly built (through {@link
|
||||
* ConnectionIdentity#panes()}, added for exactly this) and drives it with a chosen pid directly,
|
||||
* bypassing the OS-dependent PID lookup entirely — a legitimate substitute, since the pid lookup
|
||||
* is not what CB-185 is about.
|
||||
*/
|
||||
class FleetdAssemblyConnectionIdentityTest {
|
||||
|
||||
/** Same shape as {@code FleetdAssemblyLifecycleTest}'s fake, but keys {@code connectHerdr} by
|
||||
* socket path so the lead and member daemons can be two DIFFERENT {@link FakeHerdr}s. */
|
||||
private static final class TwoHerdrResourcePorts implements ResourcePorts {
|
||||
|
||||
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
|
||||
final CopyOnWriteArrayList<ScheduledExecutorService> schedulers = new CopyOnWriteArrayList<>();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
HerdrClient client = herdrsBySocket.get(socketPath);
|
||||
if (client == null) {
|
||||
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> new dev.ltms.fleet.msg.InMemoryReplyInbox();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
schedulers.add(scheduler);
|
||||
return scheduler;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never bind — this test never issues a real HTTP request.
|
||||
}
|
||||
}
|
||||
|
||||
private FleetdRuntime runtime;
|
||||
private TwoHerdrResourcePorts ports;
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
if (ports != null && ports.shutdownHook != null) {
|
||||
ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
|
||||
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
|
||||
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: "%s"
|
||||
memberHerdrSocket: "%s"
|
||||
lifecycle:
|
||||
idleTtlSeconds: 600
|
||||
health:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
""".formatted(LEAD_SOCKET, MEMBER_SOCKET));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
private FleetdRuntime assemble(Path dir, FakeHerdr lead, FakeHerdr member) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ports = new TwoHerdrResourcePorts();
|
||||
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
|
||||
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
|
||||
runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
return runtime;
|
||||
}
|
||||
|
||||
/**
|
||||
* The pin. {@code lead} carries the one pane {@link FakeHerdr}'s canned {@code
|
||||
* pane.process_info} ties to {@link FakeHerdr#WORKER_PID} (pane {@code w2:p7}); {@code member}
|
||||
* reports NO panes at all ({@link FakeHerdr#withNoPanes()}) — modelling a second daemon that
|
||||
* simply does not host the caller's pane, exactly the CB-185 javadoc's scenario for a lead's
|
||||
* own connection. If {@code PaneLocator} only ever searches the member daemon (the bug), this
|
||||
* pid resolves to nothing, because the pane that owns it lives on the LEAD daemon the bug
|
||||
* skips.
|
||||
*/
|
||||
@Test
|
||||
void connectionIdentitySearchesTheLeadDaemonNotJustTheMemberOne(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr().withNoPanes();
|
||||
|
||||
assemble(dir, lead, member);
|
||||
|
||||
PaneLocator panes = runtime.mcp().identity().panes();
|
||||
PaneLocator.Lookup lookup = panes.terminalForPid(FakeHerdr.WORKER_PID);
|
||||
|
||||
assertEquals("term_a", lookup.terminal(),
|
||||
"the pane owning WORKER_PID lives on the LEAD daemon only (the member fake reports "
|
||||
+ "no panes) — PaneLocator must still find it, which is only possible if it "
|
||||
+ "searches the lead client and not just the member one");
|
||||
}
|
||||
|
||||
/**
|
||||
* The mirror control: when the pane instead lives ONLY on the member daemon (the lead reports
|
||||
* no panes), the lookup must still find it — proving the member client is genuinely searched
|
||||
* too, not merely tolerated as a second, always-losing argument.
|
||||
*/
|
||||
@Test
|
||||
void connectionIdentityAlsoSearchesTheMemberDaemon(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr().withNoPanes();
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
|
||||
assemble(dir, lead, member);
|
||||
|
||||
PaneLocator panes = runtime.mcp().identity().panes();
|
||||
PaneLocator.Lookup lookup = panes.terminalForPid(FakeHerdr.WORKER_PID);
|
||||
|
||||
assertEquals("term_a", lookup.terminal(),
|
||||
"the pane owning WORKER_PID lives on the MEMBER daemon only — PaneLocator must "
|
||||
+ "find it there too");
|
||||
}
|
||||
|
||||
/** Sanity control: a pid nobody owns resolves to nothing on either daemon. */
|
||||
@Test
|
||||
void aPidNoPaneOwnsResolvesToNoTerminalOnEitherDaemon(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
|
||||
assemble(dir, lead, member);
|
||||
|
||||
PaneLocator panes = runtime.mcp().identity().panes();
|
||||
PaneLocator.Lookup lookup = panes.terminalForPid(999_999L);
|
||||
|
||||
assertNull(lookup.terminal());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,211 @@
|
||||
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.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.LeadChannel;
|
||||
import dev.ltms.fleet.msg.LeadChannelHandle;
|
||||
import dev.ltms.fleet.msg.LeadMessage;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
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.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 A-gaps (gap 1): {@code FleetdAssemblyLifecycleTest}'s own class javadoc says plainly
|
||||
* that it leaves {@code coordinator:} unset, so {@code leadMailbox} and {@code leadCoordLoop} stay
|
||||
* {@code null} throughout — the configured-coordinator path is never exercised by Unit A's own
|
||||
* test. This class drives that path instead: a real {@code coordinator:} block, a fake {@link
|
||||
* Fleetd.LeadMailboxOpener} returning a fake closeable channel (never a real broker connection),
|
||||
* and proof that {@link FleetdAssembly#assembleAndStart} both builds it and, on shutdown, closes it.
|
||||
*
|
||||
* <p>Made possible by generalising {@code Fleetd.LeadMailboxOpener}'s return type (and {@code
|
||||
* FleetdRuntime}'s field) from the concrete {@code LeadMailbox} to {@link LeadChannelHandle} — a
|
||||
* {@link LeadChannel} its owner can also close. {@code FleetMcp} and {@code LeadCoordLoop} already
|
||||
* consumed the narrower {@link LeadChannel}; this only widens the one seam that owns and closes it.
|
||||
*/
|
||||
class FleetdAssemblyCoordinatorLifecycleTest {
|
||||
|
||||
private static final String SELF_COORD_ID = "test-lead";
|
||||
|
||||
/** A fake {@link LeadChannelHandle}: never touches a broker, and records whether it was closed. */
|
||||
private static final class FakeLeadChannel implements LeadChannelHandle {
|
||||
volatile boolean closed = false;
|
||||
|
||||
@Override
|
||||
public void publish(String toCoordId, LeadMessage m) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<LeadMessage> peek() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String msgId) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String selfCoordId() {
|
||||
return SELF_COORD_ID;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean heldDurable() {
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public MailboxState inspect(String coordId) {
|
||||
return MailboxState.unknown(coordId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
closed = true;
|
||||
}
|
||||
}
|
||||
|
||||
/** Minimal fake {@link ResourcePorts}: a real herdr fake, a fake reply inbox, and a real, offered fake lead channel. */
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
final FakeLeadChannel leadChannel = new FakeLeadChannel();
|
||||
String offeredUri;
|
||||
String offeredSelfId;
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
offeredUri = uri;
|
||||
offeredSelfId = selfCoordId;
|
||||
return leadChannel;
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// No real HTTP bind in a unit test.
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
coordinator:
|
||||
uri: "amqp://fake-lead-broker/vh"
|
||||
selfId: "test-lead"
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@Test
|
||||
void configuredCoordinatorIsBuiltByTheAssemblyAndClosedOnShutdown(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
// --- the assembly actually calls the configured opener and builds the coordinator path ---
|
||||
assertEquals("amqp://fake-lead-broker/vh", ports.offeredUri,
|
||||
"the assembly must open the mailbox at the configured broker uri");
|
||||
assertEquals(SELF_COORD_ID, ports.offeredSelfId,
|
||||
"the assembly must open the mailbox under the configured selfId");
|
||||
assertSame(ports.leadChannel, runtime.leadMailbox(),
|
||||
"FleetdRuntime must own the exact LeadChannelHandle the opener returned, not a copy");
|
||||
assertNotNull(runtime.leadCoordLoop(),
|
||||
"a configured coordinator: block must build the receiving LeadCoordLoop too");
|
||||
assertFalse(ports.leadChannel.closed, "the channel must still be open while the daemon is running");
|
||||
|
||||
// --- shutting the assembly down closes it -------------------------------------------------
|
||||
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
|
||||
ports.shutdownHook.run();
|
||||
|
||||
assertTrue(ports.leadChannel.closed,
|
||||
"FleetdRuntime.close() must close the configured LeadChannelHandle");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,237 @@
|
||||
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.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.net.URI;
|
||||
import java.net.http.HttpClient;
|
||||
import java.net.http.HttpRequest;
|
||||
import java.net.http.HttpResponse;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 2, unit B2 (CB-185, {@code FleetApp} half). Replaces the deleted {@code
|
||||
* FleetdFleetAppConstructionTest}, which pinned this claim by reading {@code Fleetd.java}'s
|
||||
* source text for {@code "new FleetApp(herdr, memberHerdr, workers,"}. That claim moved to {@code
|
||||
* FleetdAssembly.java} (fleetd #612 Unit A) and is pinned here instead, by driving the real {@code
|
||||
* Javalin} app — via {@code runtime.app()}, not a copy — that {@link
|
||||
* FleetdAssembly#assembleAndStart} built and handed to {@link FleetdRuntime}.
|
||||
*
|
||||
* <p><strong>What this guards against</strong> (from the deleted test's own javadoc): constructing
|
||||
* {@code FleetApp} with the lead-only {@code herdr} client (dropping {@code memberHerdr}) makes
|
||||
* {@code GET /healthz} report green while the MEMBER daemon is down — so every spawn fails
|
||||
* invisibly — and silently drops every member workspace from {@code GET /sessions}. {@code
|
||||
* FleetAppTwoDaemonTest} already proves {@code FleetApp} itself merges/gates correctly given two
|
||||
* clients; the gap this pins is that the assembly actually passes it two.
|
||||
*
|
||||
* <p><strong>Both directions, not just one</strong> (fleetd #612 issue comment 17525): the deleted
|
||||
* guard's positive assertion required the exact pair {@code "new FleetApp(herdr, memberHerdr,
|
||||
* workers,"}, which does not survive EITHER daemon being dropped. An earlier version of this class
|
||||
* only proved the member-dropped direction, which left {@code new FleetApp(memberHerdr,
|
||||
* memberHerdr, ...)} — the symmetric bug, {@code /healthz} green while the LEAD daemon is down —
|
||||
* an undetected regression. {@link #healthzGoesRedWhenTheLeadDaemonIsDownEvenThoughTheMemberIsUp}
|
||||
* closes that.
|
||||
*
|
||||
* <p>Unlike the {@code ConnectionIdentity} half of CB-185 ({@code
|
||||
* FleetdAssemblyConnectionIdentityTest}), {@code /healthz} needs no caller identity at all, so
|
||||
* this test can bind {@link FleetdRuntime#app()} to a REAL ephemeral port (exactly {@code
|
||||
* FleetAppTwoDaemonTest} does for its own hand-built {@code FleetApp}) and drive it with a real
|
||||
* {@code HttpClient} — no accessor needed for this half.
|
||||
*
|
||||
* <p><strong>{@code GET /sessions} could not be driven the same way</strong>, so this class does
|
||||
* not pin the merge half of the deleted test's javadoc. {@code /sessions} requires
|
||||
* {@code Authz.Action.READ}, which — through the REAL assembly's real {@code
|
||||
* CallerResolver}/{@code ConnectionIdentity} (built with a hardcoded {@code
|
||||
* new LsofPeerPidLookup()}) — needs {@code Caller.resolved()}, i.e. a real positive pid from
|
||||
* {@code lsof}. {@code LsofPeerPidLookup} excludes its own pid (see its javadoc), and a JUnit
|
||||
* test's HTTP client and the daemon under test share one JVM pid, so the resolved pid is always
|
||||
* {@code -1} and every such request is refused as {@code ANONYMOUS} (fleetd #317's fail-closed
|
||||
* rule) before the route handler — and its {@code memberHerdr} merge — is ever reached. Verified
|
||||
* directly: driving {@code GET /sessions} here returns {@code 401 unauthenticated}, not the
|
||||
* merged body. {@code FleetAppTwoDaemonTest} avoids this because it builds {@code FleetApp} with
|
||||
* {@code callers: null}, which is not what the real assembly passes. The {@code /healthz} pin
|
||||
* below is what this class relies on for CB-185's {@code FleetApp} half; {@code
|
||||
* FleetAppTwoDaemonTest} remains the full behavioural proof that {@code FleetApp} itself merges
|
||||
* {@code /sessions} correctly once handed two clients.
|
||||
*/
|
||||
class FleetdAssemblyFleetAppTest {
|
||||
|
||||
private static final class TwoHerdrResourcePorts implements ResourcePorts {
|
||||
|
||||
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
|
||||
final CopyOnWriteArrayList<ScheduledExecutorService> schedulers = new CopyOnWriteArrayList<>();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
HerdrClient client = herdrsBySocket.get(socketPath);
|
||||
if (client == null) {
|
||||
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> new dev.ltms.fleet.msg.InMemoryReplyInbox();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
schedulers.add(scheduler);
|
||||
return scheduler;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never bind here — this test binds runtime.app() itself, for real, below.
|
||||
}
|
||||
}
|
||||
|
||||
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
|
||||
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
|
||||
|
||||
private final HttpClient http = HttpClient.newHttpClient();
|
||||
private FleetdRuntime runtime;
|
||||
private TwoHerdrResourcePorts ports;
|
||||
private Javalin boundApp;
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
if (boundApp != null) {
|
||||
boundApp.stop();
|
||||
}
|
||||
if (ports != null && ports.shutdownHook != null) {
|
||||
ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: "%s"
|
||||
memberHerdrSocket: "%s"
|
||||
lifecycle:
|
||||
idleTtlSeconds: 600
|
||||
health:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
""".formatted(LEAD_SOCKET, MEMBER_SOCKET));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
/** Assembles the real graph, then binds the real {@code Javalin app} to an ephemeral port. */
|
||||
private int assembleAndBind(Path dir, FakeHerdr lead, FakeHerdr member) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ports = new TwoHerdrResourcePorts();
|
||||
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
|
||||
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
|
||||
runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
boundApp = runtime.app().start("127.0.0.1", 0);
|
||||
return boundApp.port();
|
||||
}
|
||||
|
||||
private HttpResponse<String> get(int port, String path) throws Exception {
|
||||
HttpRequest req = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + path)).GET().build();
|
||||
return http.send(req, HttpResponse.BodyHandlers.ofString());
|
||||
}
|
||||
|
||||
/**
|
||||
* The pin. The MEMBER daemon is down; the LEAD daemon is healthy. If the assembly built
|
||||
* {@code FleetApp} with only the lead client (the bug: passing {@code herdr} where {@code
|
||||
* memberHerdr} is expected), the down member is invisible and {@code /healthz} stays 200.
|
||||
*/
|
||||
@Test
|
||||
void healthzGoesRedWhenTheMemberDaemonIsDownEvenThoughTheLeadIsUp(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr().healthy(false);
|
||||
|
||||
int port = assembleAndBind(dir, lead, member);
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
assertEquals(503, res.statusCode(),
|
||||
"a down MEMBER daemon must not be masked by a healthy lead: " + res.body());
|
||||
}
|
||||
|
||||
/** Sanity control: both daemons healthy must still be green through the real assembly. */
|
||||
@Test
|
||||
void healthzIsGreenWhenBothDaemonsAreUp(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
|
||||
int port = assembleAndBind(dir, lead, member);
|
||||
|
||||
assertEquals(200, get(port, "/healthz").statusCode());
|
||||
}
|
||||
|
||||
/**
|
||||
* The symmetric pin (fleetd #612 issue comment 17525): the LEAD daemon is down; the MEMBER
|
||||
* daemon is healthy. If the assembly built {@code FleetApp} with only the member client
|
||||
* (dropping {@code herdr} — the mirror of the bug above, {@code new FleetApp(memberHerdr,
|
||||
* memberHerdr, ...)}), the down LEAD is invisible and {@code /healthz} stays 200. Without this
|
||||
* case the pair above is one-directional and does not cover the deleted guard's positive
|
||||
* assertion (it required BOTH {@code herdr,} and {@code memberHerdr,} in that order).
|
||||
*/
|
||||
@Test
|
||||
void healthzGoesRedWhenTheLeadDaemonIsDownEvenThoughTheMemberIsUp(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr().healthy(false);
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
|
||||
int port = assembleAndBind(dir, lead, member);
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
assertEquals(503, res.statusCode(),
|
||||
"a down LEAD daemon must not be masked by a healthy member: " + res.body());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,258 @@
|
||||
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.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.inject.LoopWatchdog;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
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.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 Unit A: proves {@link FleetdAssembly#assembleAndStart} — not a copy of its logic —
|
||||
* against a real {@link FleetConfig}, a {@link FakeHerdr} and a fake {@link ResourcePorts}, with no
|
||||
* real herdr socket, no real broker, and no real HTTP bind.
|
||||
*
|
||||
* <p><strong>The canonical order this test asserts against was recorded from {@code Fleetd.java}
|
||||
* BEFORE any code moved</strong> (fleetd #612 Unit A's mandated order of work), by reading the
|
||||
* original {@code main}'s body and its shutdown-hook {@code Thread}:
|
||||
*
|
||||
* <p>Start order: {@code SessionReaper.start()} → {@code StatusPoller.start()} →
|
||||
* {@code LeadHeartbeatLoop.start()} (opt-in) → {@code FleetHealthMonitor.start()} (opt-in) →
|
||||
* {@code LeadCoordLoop.start()} (opt-in) → {@code ConfigWatcher.start()} (opt-in) →
|
||||
* {@code app.start()} (HTTP), always last.
|
||||
*
|
||||
* <p>Close order (from the original shutdown hook body): {@code sessions.close(drainTimeoutSeconds)}
|
||||
* → {@code poller.stop()} → {@code messages.close()} → {@code pushLoop.close()} →
|
||||
* {@code heartbeat.close()} (if present) → {@code leadCoordLoop.close()} (if present) →
|
||||
* {@code leadCoordScheduler.shutdownNow()} (if present) → {@code healthMonitor.stop()} (if present)
|
||||
* → {@code configWatcher.stop()} (if present) → {@code mcp.close()} → {@code reaper.stop()} (if
|
||||
* present) → {@code idleSleepGuard.close()} (if present) → {@code replyInbox.close()} (if
|
||||
* {@code AutoCloseable}) → {@code leadMailbox.close()} (if present) → {@code router.close()}.
|
||||
*
|
||||
* <p>This test's config deliberately leaves {@code coordinator:} unset, so {@code leadMailbox} and
|
||||
* {@code leadCoordLoop} stay {@code null} throughout — the lead-mailbox resource-ledger criterion is
|
||||
* NOT exercised here; see the class-level caveat in the implementer's hand-off. {@code
|
||||
* idleSleepGuard.enabled: false} is set for the same kind of reason: it would otherwise try to spawn
|
||||
* a real {@code caffeinate} subprocess, which is not one of the resources the ticket's acceptance
|
||||
* criteria names (scheduler/inbox/mailbox/loop/MCP server/herdr router).
|
||||
*/
|
||||
class FleetdAssemblyLifecycleTest {
|
||||
|
||||
/** The one {@link HerdrClient} both {@code herdrSocket} and {@code memberHerdrSocket} resolve to. */
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
final List<String> ledger = new CopyOnWriteArrayList<>();
|
||||
final List<ScheduledExecutorService> schedulers = new CopyOnWriteArrayList<>();
|
||||
final List<String> schedulerPurposes = new ArrayList<>();
|
||||
Runnable shutdownHook;
|
||||
Javalin startedApp;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
ledger.add("connectHerdr");
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
ledger.add("replyInboxOpener");
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
// Never invoked: this test's config has no `coordinator:` block, so
|
||||
// Fleetd.openLeadMailbox returns null before calling the opener at all.
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
ledger.add("newScheduler:" + purpose);
|
||||
schedulerPurposes.add(purpose);
|
||||
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
schedulers.add(scheduler);
|
||||
return scheduler;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
ledger.add("addShutdownHook");
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never call app.start(host, port): no real HTTP bind in a unit test.
|
||||
ledger.add("startHttp");
|
||||
this.startedApp = app;
|
||||
}
|
||||
}
|
||||
|
||||
/** A fake {@link ReplyInbox} that is also {@link AutoCloseable}, so the ledger can prove it closes. */
|
||||
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
volatile boolean closed = false;
|
||||
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
closed = true;
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
lifecycle:
|
||||
idleTtlSeconds: 600
|
||||
health:
|
||||
enabled: true
|
||||
intervalSeconds: 30
|
||||
leadHeartbeat:
|
||||
idleAfterSeconds: 600
|
||||
backoffMs: 15000
|
||||
quietNudgeCap: 5
|
||||
configReload:
|
||||
enabled: true
|
||||
intervalSeconds: 30
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@Test
|
||||
void assemblesTheRealBootGraphWithoutTouchingAnyRealSocketBrokerOrPort(@TempDir Path dir)
|
||||
throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
// --- start order: the herdr connect, then the recurring background loops, in the recorded
|
||||
// order, then the shutdown hook is registered, then (finally) HTTP "starts" -------------
|
||||
assertTrue(ports.ledger.indexOf("connectHerdr") < ports.ledger.indexOf("newScheduler:bridge-push-"),
|
||||
"herdr must connect before the push scheduler is created: " + ports.ledger);
|
||||
assertEquals(List.of("bridge-push-", "bridge-heartbeat-", "bridge-health-"), ports.schedulerPurposes,
|
||||
"the three always-created schedulers must be requested in exactly this order: "
|
||||
+ ports.schedulerPurposes);
|
||||
assertTrue(ports.ledger.indexOf("addShutdownHook") < ports.ledger.indexOf("startHttp"),
|
||||
"the shutdown hook must be registered before HTTP starts — the one statement that "
|
||||
+ "could not be reordered without changing FleetdRuntime's constructor shape, "
|
||||
+ "see FleetdAssembly's javadoc: " + ports.ledger);
|
||||
assertEquals(ports.ledger.size() - 1, ports.ledger.indexOf("startHttp"),
|
||||
"HTTP must start LAST of everything this fake observes: " + ports.ledger);
|
||||
assertNotNull(ports.startedApp, "FleetdAssembly must have built and handed off a real FleetApp");
|
||||
|
||||
// No real HTTP bind and no real herdr socket: this call returning at all, plus the ledger
|
||||
// above, is the proof — a real bind or a real UnixSocketHerdrClient.connect would have
|
||||
// thrown or hung against the sockets/ports this test never opened.
|
||||
assertNotNull(runtime.app(), "FleetdRuntime must own the same Javalin app that was built");
|
||||
assertEquals(ports.startedApp, runtime.app(), "attachApp must hand FleetdRuntime the SAME instance");
|
||||
|
||||
// --- the runtime owns the real, live objects — not a copy --------------------------------
|
||||
assertEquals(LoopWatchdog.State.RUNNING, runtime.reaper().health(),
|
||||
"SessionReaper must be running: lifecycle.idleTtlSeconds is configured");
|
||||
assertEquals(LoopWatchdog.State.RUNNING, runtime.poller().health(), "StatusPoller must be running");
|
||||
assertNotNull(runtime.heartbeat(), "leadHeartbeat: is configured, so the loop must be built and started");
|
||||
assertNotNull(runtime.healthMonitor(), "health.enabled: true, so the monitor must be built and started");
|
||||
assertNotNull(runtime.configWatcher(), "configReload.enabled: true, so the watcher must be built and started");
|
||||
// Documented gap (see class javadoc): no coordinator: block, so these stay null.
|
||||
assertEquals(null, runtime.leadCoordLoop(), "no coordinator: block — leadCoordLoop must stay unbuilt");
|
||||
assertEquals(null, runtime.leadMailbox(), "no coordinator: block — leadMailbox must stay unbuilt");
|
||||
assertEquals(runtime.replyInbox(), ports.replyInbox,
|
||||
"FleetdRuntime must own the exact ReplyInbox instance this fake's AmqpOpener returned");
|
||||
assertFalse(ports.herdr.closed, "herdr must still be open while the daemon is running");
|
||||
assertFalse(ports.replyInbox.closed, "the reply inbox must still be open while the daemon is running");
|
||||
for (ScheduledExecutorService scheduler : ports.schedulers) {
|
||||
assertFalse(scheduler.isShutdown(), "a scheduler must still be running while the daemon is up");
|
||||
}
|
||||
|
||||
// --- close order: invoke the captured shutdown-hook Runnable directly (no real JVM shutdown
|
||||
// happens in a unit test) and prove every resource this fake can observe is released --------
|
||||
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
|
||||
ports.shutdownHook.run();
|
||||
|
||||
assertEquals(LoopWatchdog.State.STOPPED, runtime.reaper().health(), "SessionReaper must stop on close");
|
||||
assertEquals(LoopWatchdog.State.STOPPED, runtime.poller().health(), "StatusPoller must stop on close");
|
||||
assertTrue(ports.herdr.closed, "router.close() must close the herdr client last");
|
||||
assertTrue(ports.replyInbox.closed, "the AutoCloseable reply inbox must be closed");
|
||||
for (int i = 0; i < ports.schedulers.size(); i++) {
|
||||
assertTrue(ports.schedulers.get(i).isShutdown(),
|
||||
"scheduler for '" + ports.schedulerPurposes.get(i) + "' must be shut down by close(): "
|
||||
+ "LeadHeartbeatLoop/FleetHealthMonitor/ReplyPushLoop each call "
|
||||
+ "scheduler.shutdownNow() on the exact instance ports.newScheduler(...) handed them");
|
||||
}
|
||||
|
||||
// Calling the captured hook a second time must never happen for a real JVM shutdown hook,
|
||||
// but nothing above should have thrown — that already proves every accessed field's close()
|
||||
// tolerated running once, in the recorded order, without an exception escaping.
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,165 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import io.javalin.Javalin;
|
||||
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.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 A-gaps (gap 2): {@code Fleetd.java:185} used to call {@code
|
||||
* reportRoleFallbackGaps(cfg)} <em>outside</em> the boundary {@code FleetdAssembly.assembleAndStart}
|
||||
* — the assembly call itself sat at line 201, after it — so nothing that drives the assembly (the
|
||||
* seam {@code FleetdAssemblyLifecycleTest} exercises) could ever notice the call being deleted.
|
||||
* {@code Fleetd.main} still refuses a bad config end to end (see {@code
|
||||
* FleetdStartupValidationTest}), but that only pins {@code validateAll()} and {@code
|
||||
* assertChartersNameOnlyRegisteredTools} — a validator that <em>throws</em>. {@code
|
||||
* reportRoleFallbackGaps} only logs; nothing about {@code main} throwing or not throwing can
|
||||
* observe whether that particular call ran.
|
||||
*
|
||||
* <p>This drives {@link FleetdAssembly#assembleAndStart} directly — never a copy of its logic —
|
||||
* with a config that has no {@code fleet:} pools or charters configured for any role, so every role
|
||||
* trips both of {@code reportRoleFallbackGaps}' log branches, and asserts on the real log line a
|
||||
* {@code ListAppender} attached to the shared {@code Fleetd}/{@code FleetdAssembly} logger
|
||||
* captures. Deleting the call from {@code FleetdAssembly} (verified by hand, see the ticket) turns
|
||||
* this test red; deleting it from {@code Fleetd.main} instead (its old location) would not, which
|
||||
* is exactly the gap this test closes.
|
||||
*/
|
||||
class FleetdAssemblyRoleFallbackBoundaryTest {
|
||||
|
||||
/** Minimal fake {@link ResourcePorts}: enough for {@code assembleAndStart} to run with no real I/O. */
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// No real HTTP bind in a unit test.
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
// Deliberately no `fleet:` block at all: every MemberRole has neither a pool nor a
|
||||
// charter, so reportRoleFallbackGaps' "role fallback: no fleet.<role>s: pool for ..."
|
||||
// branch is guaranteed to log something to assert on.
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@Test
|
||||
void assembleAndStartItselfReportsTheRoleFallbackGaps(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
try (var captured = CapturedLog.at(Fleetd.class, Level.INFO)) {
|
||||
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
String infoLines = captured.events().stream()
|
||||
.filter(e -> e.getLevel() == Level.INFO)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.reduce("", (a, b) -> a + "\n" + b);
|
||||
assertTrue(infoLines.contains("role fallback: no fleet.<role>s: pool for"),
|
||||
() -> "FleetdAssembly.assembleAndStart itself must call reportRoleFallbackGaps "
|
||||
+ "(fleetd #612 A-gaps gap 2) — captured INFO lines: " + infoLines);
|
||||
} finally {
|
||||
if (ports.shutdownHook != null) {
|
||||
ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,199 @@
|
||||
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.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import io.javalin.Javalin;
|
||||
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.Map;
|
||||
import java.util.OptionalLong;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — replaces {@code FleetdBackendQuarantineWiringTest} (fleetd #466), a source-text
|
||||
* test that scraped {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by fleetd #612
|
||||
* Unit A) for the {@code BackendQuarantine.withEscalation(...)} call, and separately asserted the
|
||||
* flat two-argument constructor's text was ABSENT. That proves the right method NAME appears in
|
||||
* source; it proves nothing about what the constructed object actually DOES.
|
||||
*
|
||||
* <p>This test instead drives the REAL {@link BackendQuarantine} the real {@link
|
||||
* FleetdAssembly#assembleAndStart} builds — reached through {@link
|
||||
* dev.ltms.fleet.mcp.FleetMcp#quarantineSource()} on the real, live {@code FleetMcp} {@code
|
||||
* FleetdRuntime} owns — and asserts the ONE behavioural difference {@code withEscalation} and the
|
||||
* flat constructor actually produce (see {@link BackendQuarantine}'s own class doc, "Mechanism"):
|
||||
* quarantining the same credential twice in a row, within one base cooldown of the first deadline,
|
||||
* must escalate the second cooldown past the first. A flat instance reports the identical cooldown
|
||||
* both times.
|
||||
*/
|
||||
class FleetdBackendQuarantineAssemblyTest {
|
||||
|
||||
/** Base cooldown used throughout — long enough that rounding never blurs the 2x escalation. */
|
||||
private static final int COOLDOWN_SECONDS = 100;
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
final AtomicLong clockNanos = new AtomicLong(0L);
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
// Controllable: the SAME LongSupplier instance BackendQuarantine.withEscalation(...) is
|
||||
// built with, so advancing clockNanos after assembly moves the quarantine tracker's own
|
||||
// clock, with no real sleep needed to observe escalation.
|
||||
return clockNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return clockNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
quarantineCooldownSeconds: %d
|
||||
""".formatted(COOLDOWN_SECONDS));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled BackendQuarantine escalates a repeated exhaustion, "
|
||||
+ "which the flat two-argument constructor can never do")
|
||||
void assembledQuarantineEscalatesOnARepeatedExhaustion(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
|
||||
|
||||
// First exhaustion, at clock=0: a fresh occurrence, blocked for exactly the base cooldown.
|
||||
quarantine.quarantine("cred-x");
|
||||
BackendQuarantine.Status first = quarantine.status("cred-x").orElseThrow(
|
||||
() -> new AssertionError("credential must be quarantined immediately after quarantine()"));
|
||||
assertEquals(1, first.repeatCount(), "the first call is repeat #1");
|
||||
assertEquals(COOLDOWN_SECONDS, first.remainingSeconds(),
|
||||
"a fresh quarantine blocks for exactly the base cooldown");
|
||||
|
||||
// Second exhaustion, arriving just after the first deadline — well within one base cooldown
|
||||
// of it, so this is a CONTINUATION of the same streak (repeat #2), not a fresh occurrence.
|
||||
long firstDeadlineNanos = COOLDOWN_SECONDS * 1_000_000_000L;
|
||||
ports.clockNanos.set(firstDeadlineNanos + 1);
|
||||
quarantine.quarantine("cred-x");
|
||||
BackendQuarantine.Status second = quarantine.status("cred-x").orElseThrow(
|
||||
() -> new AssertionError("credential must be quarantined immediately after the second "
|
||||
+ "quarantine() call"));
|
||||
assertEquals(2, second.repeatCount(), "the second call, arriving within one base cooldown of "
|
||||
+ "the first deadline, continues the streak as repeat #2");
|
||||
|
||||
// The one behavioural difference: withEscalation doubles the cooldown on repeat #2 (capped
|
||||
// well above this at 12x base), the flat two-argument constructor never grows past the base
|
||||
// cooldown no matter how many times quarantine() is called in a row.
|
||||
assertEquals(2 * COOLDOWN_SECONDS, second.remainingSeconds(),
|
||||
"withEscalation's default backoff doubles the cooldown on the second consecutive "
|
||||
+ "exhaustion — this is the exact call FleetdAssembly.java makes at the "
|
||||
+ "BackendQuarantine.withEscalation(...) call site");
|
||||
assertTrue(second.remainingSeconds() > first.remainingSeconds(),
|
||||
"the flat two-argument BackendQuarantine constructor would report the SAME remaining "
|
||||
+ "seconds both times — this inequality is what a mutation to the flat "
|
||||
+ "constructor at that call site must fail");
|
||||
|
||||
// Also confirm isQuarantined/remainingSeconds agree, exercising the accessors a real caller
|
||||
// (fleet_profiles / fleet_list, per BackendQuarantine's own class doc) actually reads.
|
||||
assertTrue(quarantine.isQuarantined("cred-x"));
|
||||
OptionalLong remaining = quarantine.remainingSeconds("cred-x");
|
||||
assertTrue(remaining.isPresent());
|
||||
assertEquals(2 * COOLDOWN_SECONDS, remaining.getAsLong());
|
||||
}
|
||||
}
|
||||
@@ -1,65 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #466 follow-up: {@code Fleetd.main} builds the daemon's one {@code BackendQuarantine}
|
||||
* from {@link dev.ltms.fleet.placement.BackendQuarantine#withEscalation(java.util.function.LongSupplier,
|
||||
* long)} — the escalating factory — rather than the plain two-argument constructor, which is still a
|
||||
* flat cooldown (kept for backward compatibility, see that class's doc). {@code
|
||||
* BackendQuarantineTest} proves {@code withEscalation} itself escalates, is ceilinged, and resets;
|
||||
* it says nothing about which one {@code main} actually calls.
|
||||
*
|
||||
* <p>Measured directly: reverting {@code main} to {@code new BackendQuarantine(System::nanoTime,
|
||||
* TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()))} — the pre-#466 flat call — compiles
|
||||
* with 0 errors and leaves the entire 1608-test suite (including every {@code BackendQuarantineTest}
|
||||
* case) green, because no other test constructs its {@code BackendQuarantine} through {@code main};
|
||||
* every one of them builds its own instance directly. That silent regression is exactly the shape
|
||||
* {@link FleetdLeadSeatWiringTest} and {@link FleetdCompletionResolverWiringTest} already guard
|
||||
* against for their own constructor arguments — this is the same class of gap for fleetd #466's
|
||||
* factory choice, following their approach.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
|
||||
* BackendQuarantine} and never runs {@code main} — a green result here proves only that the exact
|
||||
* text {@code main} calls {@code BackendQuarantine.withEscalation(...)} rather than the flat
|
||||
* constructor. It does not prove that call actually executes at startup (no test here starts the
|
||||
* daemon), and it does not prove the escalation reaches a real backend or credential — only
|
||||
* {@code BackendQuarantineTest} proves the factory's own behaviour, and only a live daemon proves
|
||||
* the wiring runs.
|
||||
*/
|
||||
class FleetdBackendQuarantineWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] main's BackendQuarantine local is still built from BackendQuarantine.withEscalation(...)")
|
||||
void mainStillWiresTheEscalatingQuarantineFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendQuarantine quarantine = BackendQuarantine.withEscalation(System::nanoTime,\n"
|
||||
+ " TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));"),
|
||||
"Fleetd.main's BackendQuarantine local must still be built from "
|
||||
+ "BackendQuarantine.withEscalation(System::nanoTime, "
|
||||
+ "TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds())). Reverting to the flat "
|
||||
+ "two-argument constructor (fleetd #466's measured regression) compiles with 0 errors "
|
||||
+ "and leaves the whole suite green, including every BackendQuarantineTest case that "
|
||||
+ "proves the escalation itself works — this source check is what must go red instead. "
|
||||
+ "A reverted daemon would go back to retrying a weekly subscription limit on every "
|
||||
+ "flat ~30-minute cooldown, about 336 times across the week.");
|
||||
|
||||
// Negative form of the same check: the pre-#466 flat call, if it ever reappears at this
|
||||
// declaration, must not be mistaken for the escalating one by a looser positive-only check.
|
||||
assertFalse(source.contains(
|
||||
"BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,\n"
|
||||
+ " TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));"),
|
||||
"main's BackendQuarantine local must never regress to the flat two-argument constructor");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,307 @@
|
||||
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.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.WorktreeRequest;
|
||||
import io.javalin.Javalin;
|
||||
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.Map;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 2, Unit B1 — replaces {@code FleetdCompletionResolverWiringTest} (deleted in
|
||||
* this same commit), whose four tests read {@code Fleetd.java}'s source text and asserted the
|
||||
* {@code CompletionResolver} construction call still named the right arguments. That proved the
|
||||
* call site's spelling, never that the assembled resolver actually behaves differently when an
|
||||
* argument is dropped.
|
||||
*
|
||||
* <p>These tests drive {@link FleetdAssembly#assembleAndStart} — the real boot composition,
|
||||
* fleetd #612 Unit A — and read {@link FleetdRuntime#completion()}: the exact {@link
|
||||
* CompletionResolver} instance the assembled daemon uses, never a copy built alongside it for the
|
||||
* test's benefit. Two behaviours are pinned, matching the ticket's own two measured mutations:
|
||||
*
|
||||
* <ul>
|
||||
* <li>the 8th constructor argument ({@code Fleetd.worktreeBranchLookup(sessions::roster)}) —
|
||||
* {@link #assembledResolverReportsTheMembersWorktreeAndBranchInAFallbackReport}; and</li>
|
||||
* <li>the 5th/6th arguments ({@code backendErrorPatterns}, {@code backendErrorSink}, both
|
||||
* assigned from {@code Fleetd}'s extracted factories rather than an inline lambda) —
|
||||
* {@link #assembledResolverClassifiesAndCoolsOffOnAConfiguredBackendErrorPattern}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Both tests bypass {@link dev.ltms.fleet.inject.StatusPoller} and drive {@link
|
||||
* CompletionResolver#onDelivered} / {@link CompletionResolver#resolveBeforePostAction} directly —
|
||||
* the same public, synchronous entry points {@code CompletionResolverTest} uses — with a
|
||||
* hand-built {@link CompletableFuture} waiter, so no real poller loop or herdr status poll is
|
||||
* needed. The pane scrape comes from {@link FakeHerdr#readText}; the elapsed-time floor
|
||||
* ({@code CompletionResolver.MIN_TURN_NANOS}) is controlled via a fake, advanceable {@link
|
||||
* ResourcePorts#nanoClock()} rather than a real sleep.
|
||||
*/
|
||||
class FleetdCompletionResolverAssemblyTest {
|
||||
|
||||
/** Same shape as {@code FleetdAssemblyLifecycleTest}'s fake, plus a nanoClock this test can advance. */
|
||||
private static final class ControllableResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr;
|
||||
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start
|
||||
Runnable shutdownHook;
|
||||
|
||||
ControllableResourcePorts(FakeHerdr herdr) {
|
||||
this.herdr = herdr;
|
||||
}
|
||||
|
||||
void advanceSeconds(long seconds) {
|
||||
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
// Never invoked: this test's config has no `broker:` block, so Fleetd.selectReplyInbox
|
||||
// returns the in-memory inbox before calling the opener at all.
|
||||
return (uri, prefetch) -> {
|
||||
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
// Never invoked: no `coordinator:` block configured either.
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never bind a real port.
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, String profilesYaml, String extraGuardHost,
|
||||
String worktreeRootYamlLine) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
%s
|
||||
profiles:
|
||||
%s
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- %s
|
||||
""".formatted(worktreeRootYamlLine == null ? "" : worktreeRootYamlLine, profilesYaml, extraGuardHost));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
private static void gitQuiet(Path cwd, String... args) throws Exception {
|
||||
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
|
||||
cmd.addAll(List.of(args));
|
||||
Process p = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true).start();
|
||||
String out = new String(p.getInputStream().readAllBytes());
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: git " + String.join(" ", args));
|
||||
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
|
||||
}
|
||||
|
||||
private static Path initRepo(Path dir) throws Exception {
|
||||
Files.createDirectories(dir);
|
||||
gitQuiet(dir, "init", "-q", "-b", "main");
|
||||
gitQuiet(dir, "config", "user.email", "test@example.invalid");
|
||||
gitQuiet(dir, "config", "user.name", "Test");
|
||||
Files.writeString(dir.resolve("README.md"), "seed\n");
|
||||
gitQuiet(dir, "add", "README.md");
|
||||
gitQuiet(dir, "commit", "-q", "-m", "seed");
|
||||
return dir;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248's first measured mutation: replacing {@code CompletionResolver}'s 8th constructor
|
||||
* argument with the inert {@code _ -> null} compiles clean and leaves every existing test green
|
||||
* — it silently drops fleetd #241's fallback-report location. This drives the real assembled
|
||||
* resolver through a member echoing its own injected brief back (no {@code fleet_reply}), which
|
||||
* resolves via {@code noReportMessage(target)}, and proves the real member's {@code branch} —
|
||||
* only obtainable via {@code Fleetd.worktreeBranchLookup(sessions::roster)} reading the real,
|
||||
* worktree-provisioned {@link MemberSession} — appears in the reported text.
|
||||
*/
|
||||
@Test
|
||||
void assembledResolverReportsTheMembersWorktreeAndBranchInAFallbackReport(@TempDir Path dir) throws Exception {
|
||||
Path repo = initRepo(dir.resolve("repo"));
|
||||
FleetConfig cfg = writeConfig(dir, """
|
||||
wtprofile:
|
||||
baseUrl: http://wthost.local:8000
|
||||
model: sonnet
|
||||
""", "wthost.local", "worktreeRoot: " + dir.resolve("wts"));
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
try {
|
||||
MemberSession session = runtime.sessions().acquire("wtprofile", repo.toString(), repo.toString(),
|
||||
null, new WorktreeRequest("fleetd-612-b1", null));
|
||||
String target = session.terminalId();
|
||||
String branch = session.branch();
|
||||
assertTrue(branch != null && branch.startsWith("worker/"),
|
||||
"sanity: a worktree-provisioned session must carry a real branch, got: " + branch);
|
||||
|
||||
CompletionResolver completion = runtime.completion();
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
|
||||
String echoedBrief = "z".repeat(450); // >= CompletionResolver.ECHO_MIN_CHARS normalised chars
|
||||
|
||||
ports.herdr.readText("idle, nothing yet");
|
||||
completion.onDelivered(target, new TurnToken(target, waiter, echoedBrief));
|
||||
|
||||
ports.herdr.readText(echoedBrief); // the pane just echoes the injected brief back — no real report
|
||||
ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s) without a real sleep
|
||||
completion.resolveBeforePostAction(target);
|
||||
|
||||
Rendezvous.Resolution resolution = waiter.getNow(null);
|
||||
assertTrue(resolution != null, "the waiter must have resolved synchronously");
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, resolution.kind());
|
||||
assertTrue(resolution.text().contains(CompletionResolver.NO_REPORT_PREFIX),
|
||||
"sanity: must have gone down the noReportMessage sub-path: " + resolution.text());
|
||||
assertTrue(resolution.text().contains("branch=" + branch),
|
||||
"the assembled resolver must report the member's real branch (fleetd #241 via "
|
||||
+ "fleetd #248's worktreeBranchLookup wiring); got: " + resolution.text());
|
||||
} finally {
|
||||
if (ports.shutdownHook != null) ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248's second measured mutation, and fleetd#201 Unit 5's own gap: replacing {@code
|
||||
* backendErrorPatterns}/{@code backendErrorSink} with {@code BackendErrorPatternLookup.legacy()}
|
||||
* / {@code BackendErrorSink.none()} compiles clean and leaves every existing behavioural test
|
||||
* green.
|
||||
*
|
||||
* <p>Classification proof: this test's profile configures {@code errorPattern: "credential
|
||||
* outage"} — text the built-in {@code (?i)\bAPI Error\s*:} fallback ({@code legacy()}'s only
|
||||
* behaviour) never matches. So a real {@code Fleetd.backendErrorPatternLookup(...)} wiring
|
||||
* classifies the send as {@code FAILED}; {@code legacy()} would fall through to the plain
|
||||
* completion path instead ({@code Kind.COMPLETION}).
|
||||
*
|
||||
* <p>Cool-off proof: two distinct targets on the same profile/credential each classified as a
|
||||
* backend error inside the 60s window must cool the credential off ({@link
|
||||
* dev.ltms.fleet.placement.BackendOutagePolicy}, fleetd#201 Unit 5) — observable two ways: (1)
|
||||
* the real {@code Fleetd.backendErrorSink(...)} marks each session {@code BACKEND_ERROR} (only
|
||||
* the real sink calls {@code sessions.onBackendError}; {@code BackendErrorSink.none()} never
|
||||
* does), and (2) a third explicit-profile spawn attempt is refused with a {@link
|
||||
* PlacementException} naming the cool-off — only reachable because the real sink's {@code
|
||||
* outagePolicy.record(...)} call actually ran.
|
||||
*/
|
||||
@Test
|
||||
void assembledResolverClassifiesAndCoolsOffOnAConfiguredBackendErrorPattern(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir, """
|
||||
coolprofile:
|
||||
baseUrl: http://coolhost.local:8000
|
||||
model: sonnet
|
||||
errorPattern: "credential outage"
|
||||
""", "coolhost.local", null);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
try {
|
||||
MemberSession session1 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null);
|
||||
MemberSession session2 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null);
|
||||
String target1 = session1.terminalId();
|
||||
String target2 = session2.terminalId();
|
||||
assertTrue(!target1.equals(target2), "sanity: the two spawns must be distinct targets");
|
||||
|
||||
CompletionResolver completion = runtime.completion();
|
||||
|
||||
CompletableFuture<Rendezvous.Resolution> waiter1 = new CompletableFuture<>();
|
||||
ports.herdr.readText("idle 1");
|
||||
completion.onDelivered(target1, new TurnToken(target1, waiter1, null));
|
||||
ports.herdr.readText("credential outage: upstream 503");
|
||||
ports.advanceSeconds(3);
|
||||
completion.resolveBeforePostAction(target1);
|
||||
Rendezvous.Resolution resolution1 = waiter1.getNow(null);
|
||||
assertTrue(resolution1 != null, "target1's waiter must have resolved synchronously");
|
||||
assertEquals(Rendezvous.Kind.FAILED, resolution1.kind(),
|
||||
"a configured errorPattern the built-in fallback never matches must classify as "
|
||||
+ "a backend error, not a plain completion; got: " + resolution1);
|
||||
assertTrue(resolution1.text().contains("credential outage: upstream 503"), resolution1.text());
|
||||
|
||||
CompletableFuture<Rendezvous.Resolution> waiter2 = new CompletableFuture<>();
|
||||
ports.herdr.readText("idle 2");
|
||||
completion.onDelivered(target2, new TurnToken(target2, waiter2, null));
|
||||
ports.herdr.readText("credential outage: upstream 503 again");
|
||||
ports.advanceSeconds(3);
|
||||
completion.resolveBeforePostAction(target2);
|
||||
Rendezvous.Resolution resolution2 = waiter2.getNow(null);
|
||||
assertTrue(resolution2 != null, "target2's waiter must have resolved synchronously");
|
||||
assertEquals(Rendezvous.Kind.FAILED, resolution2.kind());
|
||||
|
||||
List<MemberSession> roster = runtime.sessions().roster();
|
||||
assertTrue(roster.stream().anyMatch(s -> target1.equals(s.terminalId())
|
||||
&& s.state() == MemberSession.State.BACKEND_ERROR),
|
||||
"the real backendErrorSink must have transitioned target1 to BACKEND_ERROR: " + roster);
|
||||
assertTrue(roster.stream().anyMatch(s -> target2.equals(s.terminalId())
|
||||
&& s.state() == MemberSession.State.BACKEND_ERROR),
|
||||
"the real backendErrorSink must have transitioned target2 to BACKEND_ERROR: " + roster);
|
||||
|
||||
PlacementException coolOff = assertThrows(PlacementException.class,
|
||||
() -> runtime.sessions().acquire("coolprofile", null, dir.toString(), null),
|
||||
"two distinct targets classified within the 60s window must cool the credential "
|
||||
+ "off (BackendOutagePolicy), refusing a third explicit-profile spawn");
|
||||
assertTrue(coolOff.getMessage().contains("cooling off"), coolOff.getMessage());
|
||||
} finally {
|
||||
if (ports.shutdownHook != null) ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,89 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #248: this is the test that was actually missing. {@code Fleetd.main} builds its {@code
|
||||
* CompletionResolver} from an 8-argument constructor, and the ticket's own measurement proved two
|
||||
* ways to silently unwire it — both compiled with 0 errors and left every existing test green:
|
||||
*
|
||||
* <ul>
|
||||
* <li>replacing the worktree/branch argument (the 8th) with {@code _ -> null} — drops
|
||||
* fleetd#241's fallback-report location entirely;</li>
|
||||
* <li>replacing {@code backendErrorPatterns, backendErrorSink} (5th/6th) with {@code
|
||||
* BackendErrorPatternLookup.legacy(), BackendErrorSink.none()} — drops fleetd#201 Unit 5's
|
||||
* backend-error classification and cool-off entirely.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Neither mutation could be caught by any test that constructs its own {@code
|
||||
* CompletionResolver} (every test before this one did exactly that) or by a test of {@link
|
||||
* Fleetd#worktreeBranchLookup}, {@link Fleetd#backendErrorPatternLookup}, or {@link
|
||||
* Fleetd#backendErrorSink} in isolation (see {@code FleetdWorktreeBranchLookupTest}, {@code
|
||||
* FleetdBackendErrorPatternLookupTest}, {@code FleetdBackendErrorSinkTest}) — those prove the
|
||||
* factories work, never that {@code main} still calls them. This class is a plain source-text
|
||||
* assertion on {@code Fleetd.java} — crude, but honest about what it checks, and it turns red the
|
||||
* instant the wiring is dropped, mirroring the same fallback shape {@link
|
||||
* FleetdFleetAppConstructionTest} already uses for a different constructor argument.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
|
||||
* CompletionResolver} and never runs {@code main}.
|
||||
*/
|
||||
class FleetdCompletionResolverWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] CompletionResolver's construction call still names backendErrorPatterns and backendErrorSink")
|
||||
void backendErrorArgumentsAreStillNamedAtTheCallSite() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,"),
|
||||
"CompletionResolver's construction call must still pass backendErrorPatterns and "
|
||||
+ "backendErrorSink as its 5th/6th arguments. Replacing them with "
|
||||
+ "BackendErrorPatternLookup.legacy()/BackendErrorSink.none() (fleetd #248's measured "
|
||||
+ "mutation) compiles with 0 errors and leaves every behavioural test green — this "
|
||||
+ "source check is what must go red instead.");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] CompletionResolver's construction call still passes worktreeBranchLookup(sessions::roster)")
|
||||
void worktreeBranchLookupIsStillPassedAtTheCallSite() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("worktreeBranchLookup(sessions::roster)"),
|
||||
"CompletionResolver's construction call must still pass worktreeBranchLookup(sessions::roster) "
|
||||
+ "as its 8th (last) argument. Replacing it with the inert `_ -> null` (fleetd #248's "
|
||||
+ "other measured mutation) compiles with 0 errors and leaves every behavioural test "
|
||||
+ "green — this source check is what must go red instead.");
|
||||
assertFalse(source.contains("System::nanoTime,\n _ -> null"),
|
||||
"the worktree/branch argument must never regress to the inert `_ -> null` literal");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] backendErrorPatterns is assigned from the extracted backendErrorPatternLookup(...) factory")
|
||||
void backendErrorPatternsComesFromTheFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,"),
|
||||
"backendErrorPatterns must be assigned from Fleetd.backendErrorPatternLookup(...), not an "
|
||||
+ "inline lambda that a source check on the CompletionResolver call alone cannot see "
|
||||
+ "through");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] backendErrorSink is assigned from the extracted backendErrorSink(...) factory")
|
||||
void backendErrorSinkComesFromTheFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),"),
|
||||
"backendErrorSink must be assigned from Fleetd.backendErrorSink(...), not an inline lambda "
|
||||
+ "that a source check on the CompletionResolver call alone cannot see through");
|
||||
}
|
||||
}
|
||||
@@ -1,31 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-185: {@code ConnectionIdentity} must resolve a caller's pane on EITHER herdr daemon (a
|
||||
* lead's MCP connection resolves against the lead daemon; a member's against the member daemon).
|
||||
* Pinning {@code PaneLocator} to {@code memberHerdr} alone — the bug this guards against — leaves
|
||||
* every lead's own connection unresolvable ({@code callerTerminal == null}) the moment
|
||||
* {@code memberHerdrSocket} names a second daemon, which breaks {@code fleet_reply}/{@code
|
||||
* fleet_ask} and {@code fleet_whoami} for a lead. A unit test on {@link
|
||||
* dev.ltms.fleet.herdr.PaneLocator} alone (see {@code PaneLocatorTest}) proves the class CAN
|
||||
* search two clients, but not that {@code Fleetd.main} actually wires it that way — hence this
|
||||
* source-level assertion, the same technique {@code FleetdHerdrControlConstructionTest} uses.
|
||||
*/
|
||||
class FleetdConnectionIdentityConstructionTest {
|
||||
@Test
|
||||
void connectionIdentitySearchesBothDaemonsNotJustTheMemberOne() throws Exception {
|
||||
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
assertFalse(source.contains("new PaneLocator(memberHerdr)"),
|
||||
"PaneLocator must not be pinned to the member daemon alone — a lead's own "
|
||||
+ "connection resolves against the LEAD daemon and would never be found");
|
||||
assertTrue(source.contains("new PaneLocator(herdr, memberHerdr)"),
|
||||
"PaneLocator must search the lead daemon first, then the member daemon");
|
||||
}
|
||||
}
|
||||
@@ -1,29 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-185: {@code FleetApp} must be constructed with BOTH herdr clients (the lead's and the
|
||||
* member's), never the raw lead-only {@code herdr}. Passing only {@code herdr} — the bug this
|
||||
* guards against — makes {@code GET /healthz} green while the member daemon is down (so every
|
||||
* spawn fails invisibly) and silently drops every member workspace from {@code GET /sessions}.
|
||||
* A behavioural test on {@code FleetApp} alone (see {@code FleetAppTwoDaemonTest}) proves the
|
||||
* class merges/gates correctly when given two clients, but not that {@code Fleetd.main} actually
|
||||
* passes it two — hence this source-level assertion, mirroring
|
||||
* {@code FleetdHerdrControlConstructionTest}.
|
||||
*/
|
||||
class FleetdFleetAppConstructionTest {
|
||||
@Test
|
||||
void fleetAppIsConstructedWithBothHerdrDaemons() throws Exception {
|
||||
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
assertFalse(source.contains("new FleetApp(herdr, workers,"),
|
||||
"FleetApp must not be constructed with the lead-only herdr client");
|
||||
assertTrue(source.contains("new FleetApp(herdr, memberHerdr, workers,"),
|
||||
"FleetApp must be constructed with both the lead and the member herdr client");
|
||||
}
|
||||
}
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet;
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.msg.LeadChannelHandle;
|
||||
import dev.ltms.fleet.msg.LeadMailbox;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -35,7 +36,7 @@ class FleetdLeadMailboxSelectionTest {
|
||||
boolean unreachable;
|
||||
|
||||
@Override
|
||||
public LeadMailbox open(String uri, String selfCoordId, int prefetch) {
|
||||
public LeadChannelHandle open(String uri, String selfCoordId, int prefetch) {
|
||||
this.offeredUri = uri;
|
||||
this.offeredSelfId = selfCoordId;
|
||||
this.offeredPrefetch = prefetch;
|
||||
|
||||
@@ -0,0 +1,285 @@
|
||||
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.HerdrClient;
|
||||
import dev.ltms.fleet.lead.LeadRollover;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
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.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
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.fail;
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — replaces {@code FleetdLeadRolloverWiringTest} (fleetd #480). That class was a
|
||||
* source-text test scraping {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by
|
||||
* fleetd #612 Unit A) with three methods: {@code unrelatedAnchorStillPresent} (a scaffold anchor,
|
||||
* not an independent claim — needs no replacement of its own), {@code
|
||||
* mainStillCallsTheLeadRolloverFactory} (the call-site pin replaced by {@link
|
||||
* #assembledLeadRolloverRunsTheRealClearAndBootstrapSequence}), and {@code
|
||||
* factoryGatesOnConfigPresence} (the absent-config claim replaced by {@link
|
||||
* #absentLeadRolloverConfigMeansNoRolloverIsBuilt} — a claim this ticket found was NOT actually
|
||||
* covered behaviourally anywhere else: {@code LeadRolloverTest}'s only related assertion is
|
||||
* vacuous, {@code assertNull(null)}, and never calls the real factory).
|
||||
*
|
||||
* <p><strong>fleetd #612 B3 correction (ticket comment 17553):</strong> the first version of this
|
||||
* test configured a single shared {@link FakeHerdr} for both the lead and member herdr sockets.
|
||||
* {@code FleetdAssembly.java:140-142} falls back to {@code memberHerdr = herdr} whenever no
|
||||
* distinct {@code memberHerdrSocket} is configured, so with one fake, {@code
|
||||
* router.leadAgents()} and {@code router.memberAgents()} wrapped the identical client — a
|
||||
* mutation swapping {@code Fleetd.leadRollover(cfg, router.leadAgents(), config, leads)} for
|
||||
* {@code ..., router.memberAgents(), ...} at {@code FleetdAssembly.java:408} was therefore
|
||||
* invisible to this test, even though the two are genuinely different daemons in production. This
|
||||
* version configures two distinct sockets and two distinct {@link FakeHerdr} instances (the same
|
||||
* pattern {@code FleetdAssemblyConnectionIdentityTest}, fleetd #612 B2, already uses to separate
|
||||
* lead from member) and asserts the roll's {@code /clear}/bootstrap sends land on the LEAD fake
|
||||
* and never on the MEMBER one.
|
||||
*/
|
||||
class FleetdLeadRolloverAssemblyTest {
|
||||
|
||||
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
|
||||
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
|
||||
|
||||
/** Keys {@code connectHerdr} by socket path so the lead and member daemons can be two
|
||||
* DIFFERENT {@link FakeHerdr}s — same shape as B2's {@code FleetdAssemblyConnectionIdentityTest
|
||||
* .TwoHerdrResourcePorts}. */
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
HerdrClient client = herdrsBySocket.get(socketPath);
|
||||
if (client == null) {
|
||||
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, Path leadCwd) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: "%s"
|
||||
memberHerdrSocket: "%s"
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus"
|
||||
cwd: "%s"
|
||||
leadRollover:
|
||||
handoverPath: handover.md
|
||||
requireOperatorConfirm: false
|
||||
""".formatted(LEAD_SOCKET, MEMBER_SOCKET, leadCwd.toString()));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled LeadRollover runs the full open/confirm/continuation "
|
||||
+ "sequence — /clear, then bootstrapText — through the real herdr router")
|
||||
void assembledLeadRolloverRunsTheRealClearAndBootstrapSequence(@TempDir Path dir) throws Exception {
|
||||
Path leadCwd = dir.resolve("lead-workspace");
|
||||
Files.createDirectories(leadCwd);
|
||||
FleetConfig cfg = writeConfig(dir, leadCwd);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
// Two DISTINCT fakes — one per configured socket — so leadAgents()/memberAgents() wrap
|
||||
// genuinely different clients, exactly like production when memberHerdrSocket is set.
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
lead.withTab("w2", "w2:t7", "lead: opus");
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
|
||||
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
LeadRollover rollover = runtime.mcp().leadRollover();
|
||||
assertNotNull(rollover, "leadRollover: is present in this test's config, so "
|
||||
+ "FleetdAssembly.assembleAndStart must have built a real LeadRollover through the "
|
||||
+ "Fleetd.leadRollover(...) call site — a mutation to `LeadRollover leadRollover = "
|
||||
+ "null;` at that call site can never pass this");
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open("term_a", "fleetd #612 B3 test");
|
||||
String expectedHandoverPath = leadCwd.resolve("handover.md").normalize().toString();
|
||||
assertEquals(expectedHandoverPath, pending.handoverPath());
|
||||
|
||||
// Ensure the handover file's mtime lands strictly AFTER open()'s requestedAtMillis —
|
||||
// LeadRollover.checkHandover refuses on mtime <= requestedAt (HANDOVER_STALE).
|
||||
Thread.sleep(50);
|
||||
Files.writeString(Path.of(pending.handoverPath()), "handover content for fleetd #612 B3");
|
||||
|
||||
LeadRollover.RollDecision decision = rollover.confirm("term_a", pending.token(), true);
|
||||
assertTrue(decision.accepted(), "confirm() must approve: requireOperatorConfirm is false, "
|
||||
+ "the caller terminal matches open()'s, and the handover file exists, is non-empty "
|
||||
+ "and fresh — got: " + decision);
|
||||
|
||||
// The production LeadRollover constructor always runs the post-confirm continuation on a
|
||||
// real virtual thread (see Fleetd.leadRollover, which never passes the package-private test
|
||||
// constructor), so this polls the real FleetMcp.leadRollover() instance's status(token)
|
||||
// until the real continuation finishes.
|
||||
LeadRollover.RollStatus status = pollUntilTerminal(rollover, pending.token());
|
||||
assertEquals(LeadRollover.RollState.ROLLED, status.state(),
|
||||
"the full happy path must complete: FakeHerdr's default agent status is 'idle', so "
|
||||
+ "the turn-boundary wait settles immediately and the post-/clear wait "
|
||||
+ "releases via its pickup-grace path — detail: " + status.detail());
|
||||
|
||||
// Prove the real herdr router actually sent BOTH messages, in order, to the real LEAD
|
||||
// pane — this is the one thing a source-text pin on the call site could never show.
|
||||
List<FakeHerdr.Call> prompts = lead.calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt"))
|
||||
.toList();
|
||||
assertTrue(prompts.size() >= 2, "expected at least a /clear send and a bootstrapText send "
|
||||
+ "on the LEAD daemon, got " + prompts.size() + " agent.prompt calls: " + prompts);
|
||||
assertEquals("/clear", ((Map<String, Object>) prompts.get(0).params()).get("text"),
|
||||
"the first send must be the literal /clear housekeeping command");
|
||||
Object secondText = ((Map<String, Object>) prompts.get(1).params()).get("text");
|
||||
assertTrue(secondText instanceof String && ((String) secondText).contains(expectedHandoverPath),
|
||||
"the second send must be the default bootstrapText naming the resolved handover "
|
||||
+ "path, got: " + secondText);
|
||||
|
||||
// fleetd #612 B3 correction: prove the roll never touches the MEMBER daemon. A mutation
|
||||
// swapping router.leadAgents() for router.memberAgents() at the real call site would move
|
||||
// both sends above onto `member` instead, which this assertion catches — the thing the
|
||||
// single-fake version of this test could never see, because both wrapped the same client.
|
||||
List<FakeHerdr.Call> memberPrompts = member.calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt"))
|
||||
.toList();
|
||||
assertTrue(memberPrompts.isEmpty(), "the roll must be wired to the LEAD daemon only — got "
|
||||
+ memberPrompts.size() + " agent.prompt call(s) on the MEMBER daemon instead: "
|
||||
+ memberPrompts);
|
||||
}
|
||||
|
||||
private static LeadRollover.RollStatus pollUntilTerminal(LeadRollover rollover, String token)
|
||||
throws InterruptedException {
|
||||
long deadline = System.nanoTime() + java.util.concurrent.TimeUnit.SECONDS.toNanos(10);
|
||||
while (System.nanoTime() < deadline) {
|
||||
LeadRollover.RollStatus status = rollover.status(token);
|
||||
if (status.state() != LeadRollover.RollState.PENDING
|
||||
&& status.state() != LeadRollover.RollState.IN_PROGRESS) {
|
||||
return status;
|
||||
}
|
||||
Thread.sleep(50);
|
||||
}
|
||||
fail("the real continuation did not reach a terminal state within 10s — last status: "
|
||||
+ rollover.status(token));
|
||||
throw new AssertionError("unreachable");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] Fleetd.leadRollover(...) returns null when leadRollover: is absent "
|
||||
+ "from config — the opt-in gate FleetdLeadRolloverWiringTest's "
|
||||
+ "factoryGatesOnConfigPresence pinned by source text alone")
|
||||
void absentLeadRolloverConfigMeansNoRolloverIsBuilt(@TempDir Path dir) throws Exception {
|
||||
Path yaml = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
""");
|
||||
ConfigRef config = new ConfigRef(yaml, FleetConfig.load(yaml));
|
||||
AgentControl agents = new AgentControl(new FakeHerdr());
|
||||
|
||||
LeadRollover rollover = Fleetd.leadRollover(config.get(), agents, config, Map::of);
|
||||
|
||||
assertNull(rollover, "leadRollover: is absent from this config, so the factory's opt-in "
|
||||
+ "gate (`if (cfg.leadRollover() == null) return null;`) must fire and no "
|
||||
+ "LeadRollover must be constructed at all");
|
||||
}
|
||||
}
|
||||
@@ -1,89 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #480 Unit A, hard requirement 6: pin {@code Fleetd.main}'s construction of {@link
|
||||
* dev.ltms.fleet.lead.LeadRollover} with a source-text assertion, mirroring {@code
|
||||
* FleetdCompletionResolverWiringTest}'s pattern — five log-only reporters in {@code Fleetd.main}
|
||||
* already survived mutation batteries this exact way (fleetd #415's extraction antidote note).
|
||||
*
|
||||
* <p>What this class still covers, and what it never claimed to. {@code LeadRolloverTest}
|
||||
* constructs its own {@code LeadRollover} directly (as every prior test of an extracted factory
|
||||
* does) with a hand-built lookup, so a mutation that deletes the {@code leadRollover(...)} call
|
||||
* from {@code main} — or replaces one of its arguments with something that still compiles, e.g.
|
||||
* {@code router.leadAgents()} swapped for {@code null}, or the whole assignment swapped for a bare
|
||||
* {@code null} literal — leaves every behavioural test green. This is a plain string read, guarded
|
||||
* by an unrelated anchor assertion so a broken or empty file read cannot pass as a real change.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
|
||||
* LeadRollover} and never runs {@code main}. It pins the {@code leadRollover(...)} CALL SITE's
|
||||
* argument list — that {@code main} still passes {@code leads} at all — never what the factory
|
||||
* DOES with that argument once inside its own body.
|
||||
*
|
||||
* <p><b>Correction (fleetd #480 relative-handover-path follow-up): that gap used to be real, and
|
||||
* now is not — but not here.</b> This class's javadoc previously claimed "no behavioural test can
|
||||
* catch this wiring dropping out" for the whole factory, including the lambda {@code
|
||||
* leadRollover(...)} builds internally (terminal → lead name → {@code Leader.cwd()}). That claim
|
||||
* was proven true at the time — mutating that lambda's body to {@code String leadName = null;}
|
||||
* (always "no lead found", which silently reintroduces the daemon-cwd bug this ticket fixes) left
|
||||
* the full suite green, {@code Tests run: 1669, Failures: 0}. It is no longer true: {@code
|
||||
* FleetdLeadRolloverWorkspaceLookupTest} now calls {@code Fleetd.leadRollover(...)} directly with a
|
||||
* real {@link dev.ltms.fleet.config.ConfigRef} built from a temp {@code fleetd.yaml}, and fails
|
||||
* against that exact one-line mutation. So: THIS class still covers only the call site's argument
|
||||
* list; {@code FleetdLeadRolloverWorkspaceLookupTest} is what now covers the lambda's body. Neither
|
||||
* one subsumes the other — keep both.
|
||||
*/
|
||||
class FleetdLeadRolloverWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] unrelated anchor: Fleetd.java still declares the Fleetd class")
|
||||
void unrelatedAnchorStillPresent() throws Exception {
|
||||
// Guards the two assertions below: without this, a bad read (empty string, wrong file,
|
||||
// truncated file) could vacuously fail to contain the leadRollover(...) call too, and a
|
||||
// test that only asserts "contains X" would report a false pass for the wrong reason if X
|
||||
// happened to match. Asserting an unrelated, structurally distant string first proves the
|
||||
// read actually pulled real file content.
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("public final class Fleetd"),
|
||||
"sanity anchor failed — the file read did not return real Fleetd.java source; the "
|
||||
+ "leadRollover(...) wiring assertions below cannot be trusted until this passes");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] main still constructs LeadRollover via the leadRollover(...) factory, exactly as heartbeat is constructed")
|
||||
void mainStillCallsTheLeadRolloverFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"LeadRollover leadRollover = leadRollover(cfg, router.leadAgents(), config, leads);"),
|
||||
"Fleetd.main must still assign `LeadRollover leadRollover = leadRollover(cfg, "
|
||||
+ "router.leadAgents(), config, leads);`. Dropping this call, or swapping one of "
|
||||
+ "its arguments for something that still compiles (e.g. null in place of "
|
||||
+ "router.leadAgents()), leaves every behavioural test green — this source check is "
|
||||
+ "what must go red instead. fleetd #480 correction 2 deliberately dropped "
|
||||
+ "primaryRegistry from this call — see LeadRollover's class javadoc for why a "
|
||||
+ "single-slot lookup was wrong here. The fleetd #480 relative-handover-path "
|
||||
+ "follow-up added `leads` (terminal → lead name) so the factory can resolve a "
|
||||
+ "relative handoverPath against the calling lead's own workspace.");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] the leadRollover(...) factory itself gates construction on cfg.leadRollover() != null")
|
||||
void factoryGatesOnConfigPresence() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("if (cfg.leadRollover() == null) {"),
|
||||
"Fleetd.leadRollover(...) must refuse to construct a LeadRollover when the "
|
||||
+ "leadRollover: block is absent — an upgraded daemon must never silently acquire "
|
||||
+ "the ability to clear the lead's own pane. See LeadHeartbeatLoop's construction "
|
||||
+ "gate (cfg.leadHeartbeat() != null) for the pattern this mirrors.");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,172 @@
|
||||
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.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
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.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — replaces {@code FleetdLeadSeatWiringTest} (fleetd #176), a source-text test that
|
||||
* scraped {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by fleetd #612 Unit A)
|
||||
* for the exact {@code new FleetMcp.LeadSeatSource(Fleetd.leadSeatLookup(...))} constructor-call
|
||||
* text. That proves the right symbols appear in source; it proves nothing about what the daemon's
|
||||
* live {@code fleet_list} actually reports.
|
||||
*
|
||||
* <p>This test instead drives the REAL {@link FleetMcp.LeadSeatSource} the real {@link
|
||||
* FleetdAssembly#assembleAndStart} builds — including the REAL {@code LeadTabScanner} it wires
|
||||
* {@code Fleetd.leadSeatLookup} through — reached via {@link FleetMcp#leadSeatSource()} on the
|
||||
* live {@code FleetMcp} {@code FleetdRuntime} owns. It seeds one FakeHerdr tab labelled to match a
|
||||
* configured {@code fleet.leaders.opus.tab}, with a live agent already in it (FakeHerdr's own
|
||||
* default {@code agent.list}/{@code pane.list} entries for {@code term_a}/{@code w2:p7}/{@code
|
||||
* w2:t7} — no FakeHerdr change needed), and asserts the assembled seat source reports exactly the
|
||||
* seat {@link FleetMcp.LeadSeatSource#none()} (the inert stand-in) could never produce: 1, not 0.
|
||||
*/
|
||||
class FleetdLeadSeatAssemblyTest {
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus"
|
||||
profile: sonnet
|
||||
profiles:
|
||||
sonnet:
|
||||
subscription: true
|
||||
argv: ["ccs", "sonnet"]
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled LeadSeatSource, backed by the real LeadTabScanner, "
|
||||
+ "reports a live lead's seat against its own subscription profile")
|
||||
void assembledLeadSeatSourceReportsALiveLeadsSeat(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
// Label FakeHerdr's own default pane's tab (term_a / w2:p7 / w2:t7, already carrying a live
|
||||
// agent) to match fleet.leaders.opus.tab exactly — no FakeHerdr change needed at all.
|
||||
ports.herdr.withTab("w2", "w2:t7", "lead: opus");
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
FleetMcp.LeadSeatSource seatSource = runtime.mcp().leadSeatSource();
|
||||
assertEquals(1, seatSource.seatsFor().apply("sonnet"),
|
||||
"the real LeadTabScanner recognises the labelled tab as a live 'opus' lead on "
|
||||
+ "profile 'sonnet' (same credential, matched by Fleetd.leadSeatLookup), so "
|
||||
+ "subscription profile 'sonnet' must be charged one seat — "
|
||||
+ "FleetMcp.LeadSeatSource.none() (the inert stand-in this test's mutation "
|
||||
+ "swaps the call site for) always reports 0, whatever the input");
|
||||
|
||||
// A profile no lead is running on gets no seat charged — the same seat source, applied to
|
||||
// an input that must stay at the inert answer even on the real, non-inert instance.
|
||||
assertEquals(0, seatSource.seatsFor().apply("no-such-profile"));
|
||||
}
|
||||
}
|
||||
@@ -1,43 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #176: {@code Fleetd.main} builds its {@code FleetMcp} from a 14-argument constructor whose
|
||||
* last argument is a {@code FleetMcp.LeadSeatSource} wrapping {@link Fleetd#leadSeatLookup}. That
|
||||
* argument is exactly the kind of wiring fleetd #248 warned about: dropping it (or swapping it for
|
||||
* the inert {@code FleetMcp.LeadSeatSource.none()}) compiles with 0 errors and leaves every test
|
||||
* that builds its own {@code FleetMcp}/{@code CapacitySource} directly — every test that predates
|
||||
* this ticket — green, because none of them go through {@code main} at all.
|
||||
*
|
||||
* <p>{@link FleetdLeadSeatLookupTest} proves the factory's own matching logic; this class is the
|
||||
* plain source-text assertion that proves {@code main} still passes its result in, mirroring
|
||||
* {@code FleetdCompletionResolverWiringTest}'s approach for the same class of gap.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a
|
||||
* {@code FleetMcp} and never runs {@code main}.
|
||||
*/
|
||||
class FleetdLeadSeatWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] FleetMcp's construction call still passes a LeadSeatSource built from leadSeatLookup(...)")
|
||||
void fleetMcpConstructionStillWiresLeadSeatLookup() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), "
|
||||
+ "leaders, leads))"),
|
||||
"FleetMcp's construction call must still pass a LeadSeatSource built from "
|
||||
+ "Fleetd.leadSeatLookup(...). Dropping it or swapping in "
|
||||
+ "FleetMcp.LeadSeatSource.none() (fleetd #176's would-be silent regression, the same "
|
||||
+ "shape as fleetd #248's measured mutations) compiles with 0 errors and leaves every "
|
||||
+ "existing behavioural test green — this source check is what must go red instead.");
|
||||
}
|
||||
}
|
||||
@@ -456,7 +456,11 @@ public final class FakeHerdr implements HerdrClient {
|
||||
}
|
||||
}
|
||||
|
||||
/** fleetd #612 Unit A: set by {@link #close()} so a test can prove a router/client actually closed this. */
|
||||
public volatile boolean closed = false;
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
closed = true;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -476,6 +476,26 @@ class LeadHeartbeatLoopTest {
|
||||
assertTrue(notice.contains("2 compactions"), notice);
|
||||
}
|
||||
|
||||
// ── fleetd #621: the notice must track the effective requireOperatorConfirm value ─────────────
|
||||
|
||||
@Test
|
||||
void contextNoticeKeepsAskingTheOperatorWhenRequireOperatorConfirmIsTrue() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
String notice = LeadHeartbeatLoop.contextNotice(true, reading, false, true);
|
||||
assertTrue(notice.contains("ask the operator"), notice);
|
||||
assertTrue(notice.contains("Only the operator can approve the roll"), notice);
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeDropsTheOperatorAskWhenRequireOperatorConfirmIsFalse() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
String notice = LeadHeartbeatLoop.contextNotice(true, reading, false, false);
|
||||
assertFalse(notice.contains("ask the operator"), notice);
|
||||
assertFalse(notice.contains("Only the operator can approve the roll"), notice);
|
||||
assertTrue(notice.contains("fleet_handover"), notice);
|
||||
assertTrue(notice.contains("maxDocAgeSeconds"), notice);
|
||||
}
|
||||
|
||||
// ── fleetd #609 review: the latch must mean "the notice reached the pane" ────────────────────
|
||||
//
|
||||
// These four drive LeadHeartbeatLoop.tick() directly (package-private, same reasoning as
|
||||
|
||||
@@ -1846,7 +1846,8 @@ class MessageServiceTest {
|
||||
* but backed by {@link ManualScheduler} instead of a real timer (fleetd #608): its tick never
|
||||
* fires on its own — a test drives it explicitly via {@link ManualScheduler#runDueTasks()}.
|
||||
*/
|
||||
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler)
|
||||
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler,
|
||||
ReplyPushLoop pushLoop)
|
||||
implements AutoCloseable {
|
||||
@Override
|
||||
public void close() {
|
||||
@@ -1861,14 +1862,20 @@ class MessageServiceTest {
|
||||
* {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that.
|
||||
*/
|
||||
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) {
|
||||
return wireWithManualScheduler(maxReminders, backoffMs, System::nanoTime);
|
||||
}
|
||||
|
||||
/** As above, with an injectable clock for tests that exercise terminal-ticket pruning. */
|
||||
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs,
|
||||
java.util.function.LongSupplier nowNanos) {
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(T, LEAD);
|
||||
FakeHerdr leadHerdr = new FakeHerdr();
|
||||
AgentControl leadAgents = new AgentControl(leadHerdr);
|
||||
ManualScheduler scheduler = new ManualScheduler();
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime);
|
||||
return new ManualPushWiring(service, leadHerdr, scheduler);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, nowNanos);
|
||||
return new ManualPushWiring(service, leadHerdr, scheduler, pushLoop);
|
||||
}
|
||||
|
||||
private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException {
|
||||
@@ -1959,24 +1966,35 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception {
|
||||
try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: both tickets land before the tick fires
|
||||
String first = wiring.service().sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "first done"));
|
||||
// Settle without polling: poll() itself marks a ticket collected (that's the point of
|
||||
// anAlreadyCollectedTicketProducesNoNudge above) — using it here to detect completion
|
||||
// would collect the ticket before the coalescing this test checks ever gets a chance.
|
||||
Thread.sleep(100);
|
||||
// A 1ms backoff is due immediately. ManualScheduler still cannot run it until this test
|
||||
// explicitly calls runDueTasks(), so both terminal tickets join one scheduled tick.
|
||||
try (var wiring = wireWithManualScheduler(1, 1)) {
|
||||
String first;
|
||||
String second;
|
||||
java.util.concurrent.CountDownLatch firstTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(firstTerminalReached::countDown);
|
||||
try {
|
||||
first = wiring.service().sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "first done"));
|
||||
assertTrue(firstTerminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the first ticket never reached its terminal phase");
|
||||
|
||||
String second = wiring.service().sendAsync(T, "second task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "second done"));
|
||||
Thread.sleep(100);
|
||||
java.util.concurrent.CountDownLatch secondTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(secondTerminalReached::countDown);
|
||||
second = wiring.service().sendAsync(T, "second task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "second done"));
|
||||
assertTrue(secondTerminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the second ticket never reached its terminal phase");
|
||||
} finally {
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||
}
|
||||
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
Thread.sleep(200); // settle — nothing more should arrive beyond the one coalesced nudge
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(),
|
||||
"both terminal tickets must coalesce onto one scheduled tick");
|
||||
long nudgeCount = wiring.leadHerdr().calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||
assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two");
|
||||
@@ -2055,7 +2073,8 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception {
|
||||
try (var wiring = wireWithPushLoop(5, 50)) {
|
||||
// A 1ms backoff is due immediately, but ManualScheduler only ticks when this test asks it to.
|
||||
try (var wiring = wireWithManualScheduler(5, 1)) {
|
||||
String ticket = wiring.service().sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
@@ -2063,8 +2082,9 @@ class MessageServiceTest {
|
||||
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> wiring.service().ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
|
||||
awaitQuestionPendingOn(wiring.pushLoop(), LEAD, asking.turnId());
|
||||
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "the open question must have one scheduled tick");
|
||||
long callsBeforeAnswer = wiring.leadHerdr().calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||
|
||||
@@ -2072,12 +2092,22 @@ class MessageServiceTest {
|
||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(terminalReached::countDown);
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
try {
|
||||
assertTrue(terminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the answered ticket never reached its terminal phase");
|
||||
} finally {
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||
}
|
||||
answer.get(5, TimeUnit.SECONDS);
|
||||
|
||||
// Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire —
|
||||
// none of them may still name the question's turnId, which is closed.
|
||||
Thread.sleep(300);
|
||||
// Run the ticket's legitimate terminal nudge and every later scheduled tick through its
|
||||
// reminder cap. None may still name the closed question.
|
||||
for (int tick = 0; tick < 6; tick++) {
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "expected one scheduled reminder tick");
|
||||
}
|
||||
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt"))
|
||||
.skip(callsBeforeAnswer)
|
||||
@@ -2160,35 +2190,45 @@ class MessageServiceTest {
|
||||
// decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix)
|
||||
// pendingTickets never drops the ticket. That is the exact "cap already STOPped" branch of
|
||||
// the bug report, reached deterministically rather than by timing it against a live tick.
|
||||
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
|
||||
String stale = wiring.service().sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "stale result"));
|
||||
// A 1ms backoff is due immediately, but ManualScheduler runs only the ticks below.
|
||||
try (var wiring = wireWithManualScheduler(1, 1, clock::get)) {
|
||||
String stale;
|
||||
String fresh;
|
||||
java.util.concurrent.CountDownLatch staleTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(staleTerminalReached::countDown);
|
||||
try {
|
||||
stale = wiring.service().sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "stale result"));
|
||||
assertTrue(staleTerminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the stale ticket never reached its terminal phase");
|
||||
|
||||
// Let the reminder loop fire its one nudge and hit the cap (STOP removes it from
|
||||
// activeLeads; pendingTickets is untouched either way — that asymmetry is the bug).
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
Thread.sleep(300);
|
||||
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
|
||||
"sanity: the stale ticket's own reminder must have fired first");
|
||||
// Fire the stale ticket's one nudge, then its cap tick (STOP removes it from activeLeads;
|
||||
// pendingTickets is untouched either way — that asymmetry is the bug).
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one nudge tick");
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one cap tick");
|
||||
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
|
||||
"sanity: the stale ticket's own reminder must have fired first");
|
||||
|
||||
// Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly.
|
||||
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
|
||||
// Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly.
|
||||
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
|
||||
|
||||
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
|
||||
String fresh = wiring.service().sendAsync(T, "second task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "fresh result"));
|
||||
|
||||
// The fresh ticket restarts the (now-dormant) reminder loop with its own nudge.
|
||||
long before = wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count();
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count() <= before
|
||||
&& System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(10);
|
||||
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
|
||||
java.util.concurrent.CountDownLatch freshTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(freshTerminalReached::countDown);
|
||||
fresh = wiring.service().sendAsync(T, "second task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "fresh result"));
|
||||
assertTrue(freshTerminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the fresh ticket never reached its terminal phase");
|
||||
} finally {
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||
}
|
||||
|
||||
// The fresh ticket restarts the now-dormant reminder loop with its own nudge.
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "the fresh ticket must have one nudge tick");
|
||||
String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
||||
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
|
||||
assertFalse(latestNudge.contains(stale),
|
||||
|
||||
Reference in New Issue
Block a user