Compare commits

..

21 Commits

Author SHA1 Message Date
Dai Ha 3437d6313d fleetd #241: bound the echo match so a real report is never swallowed
CI / contract (pull_request) Successful in 1m9s
CI / build (pull_request) Successful in 1m25s
Round 1 used plain bidirectional containment. The direction that catches the real bug --
the pane holds the brief plus a status bar, so the scrape contains the brief -- also fires
when a member restates the whole brief and then writes a genuine report under it. That
threw the report away and told the lead nothing was produced, which is worse than the bug
being fixed: it destroys a delivery instead of merely obscuring one.

The safe direction (the scrape is a fragment of the brief) stays unbounded, because a
fragment of the brief is by definition not a report. The dangerous direction now requires
the scrape to add at most MAX_ECHO_EXCESS_CHARS beyond the brief, which is the amount of
TUI chrome a real echo carries.

Work by the cb241 worker, committed by the lead: its backend stopped answering after the
fix was written, so two turns ended with no commit and no reply. Verified by the lead:
1169 tests, 0 failures; removing the bound turns pinsTheMaximumTuiChromeExcess and
completionFallbackKeepsARealReportThatRestatesTheWholeBrief red with 0 compile errors.
2026-09-03 11:39:11 +07:00
Dai Ha 321d8dcbb5 fleetd #241: suppress echoed fallback briefs
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m20s
2026-09-03 11:10:48 +07:00
Dai Ha 26bafe824b fleetd #201/#227 unit 1: classify a backend error at the scrape
CI / contract (push) Successful in 43s
CI / build (push) Successful in 1m54s
CompletionResolver already reads the pane on every turn, so the classifier lives there
rather than in a new watcher. A matched pattern resolves the waiter as a failure and
fires BackendErrorSink, always inside the resolveFailure win-gate so exactly one thread
reports one incident.

The too-fast path now takes a fresh scrape instead of reusing the pre-turn text, so a
backend that dies immediately is still classified. BackendErrorSink is a single-method
functional interface by design — see #234 for what a default overload does to a lambda.
2026-09-03 10:53:51 +07:00
Dai Ha 838a701109 fleetd #234: carry the profile hint to the exhaustion sink
An exhausted opencode member could not always be mapped back to a credential, because
the launcher knew the profile and the sink did not. The sink now takes a profile hint.

The interface is inverted on purpose: the three-argument method is the single abstract
method and the two-argument one is the default. A lambda can only implement the abstract
method, so every lambda is now forced to carry the profile. The first round of this fix
added the third argument as a default overload, and the production forwarder in Fleetd
was a two-argument lambda — so the fix compiled, passed its tests, and never ran.

Verified by the lead: deleting the forwardingTo factory's override, and rewriting the
Fleetd call site as a plain lambda, both fail to compile now rather than passing silently.
2026-09-03 10:52:01 +07:00
Dai Ha e5eb3534c7 fleetd #201/#227 unit 3: lead outage nudge
Backend incidents become a fourth source inside ReplyPushLoop, not a new scheduler — a
second injector would race the one control that already owns lead-pane delivery. One
notice per incident per affected lead (not per member), one-shot, waiting while the lead
pane is not injectable, and combined into the same nudge as any pending failed ticket.

A classified target that cannot be mapped to a credential gets its own truthful notice
and its own pending/delivered records. It previously reused the incident message, which
told the lead a credential was 'cooling for 0 remaining seconds' when nothing was cooling,
and smuggled the free-text reason into the profiles field.

Verified by the lead: 1129 tests green; routing the unmapped path back through
onBackendIncident turns unmappedBackendTargetUsesTheKnownLeadSchedule red with 0 compile
errors, and that test now asserts the whole rendered message rather than two substrings.
2026-09-03 10:49:25 +07:00
Dai Ha 959c83534f fleetd #201/#227 unit 2: credential outage policy
Adds BackendOutagePolicy — a credential-keyed state machine on an injected monotonic
clock. Two classified backend errors from two DISTINCT targets on one credential inside
60 seconds mint one incident and start a 60-second cool-off. Errors during cool-off
neither extend it nor mint another; expiry clears evidence, so two fresh errors rearm.

Correlated on credentialId, never on profile name or error text. Deliberately not
BackendQuarantine: that restarts a 1800-second cooldown per exhaustion, and its name
would make every refusal say 'backend exhausted', which is a different condition.

Threshold counts distinct targets rather than raw events (lead decision): the classifier
is a heuristic and a valid member report can quote an 'API Error:' line, so one member
repeating that line must not remove a healthy credential's capacity. A real outage hits
every member on the credential, so true detection is unaffected.

Verified by the lead: 1135 tests green; reverting evidenceCount() to reasons.size()
turns two BackendOutagePolicyTest cases red with 0 compile errors.
2026-09-03 10:49:10 +07:00
Dai Ha 31b028e860 fleetd #234 round 4: invert ExhaustionSink's abstract method so the bug class is unrepresentable
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m23s
Round 3's factory fixed the two known call sites but the underlying shape
was still there: a lambda written against ExhaustionSink binds to whichever
overload is abstract, and the 2-arg form held that position, so ANY lambda
-- a call-site forwarder, a hand-built test double, a future caller who has
never heard of fleetd #234 -- could still silently take the hint-dropping
default. Two rounds shipped exactly that mistake in two different places.

Fix: made the 3-arg onExhausted(target, reason, profile) the interface's
single abstract method; the 2-arg form is now a default that delegates with
a null profile. A lambda declared against ExhaustionSink today is forced by
the compiler to take three parameters -- there is no overload left for it to
bind to that can drop the hint. This is enforced by the type system, not by
a test that has to remember to check for it.

Knock-on changes:
- ExhaustionSink.none() -- a 3-arg lambda, still a genuine no-op, now safe
  by construction rather than by care.
- ExhaustionSink.forwardingTo(...) -- collapses to a one-line 3-arg lambda;
  kept as a named factory (round 3's lesson: a test must call the real
  object, not rebuild its shape).
- Fleetd.java's real sink and the two OpenCodeLauncherTest sinks that used
  to be anonymous classes overriding both overloads are now plain lambdas
  too -- the 2-arg override each carried was pure boilerplate once the
  interface provides it as a default.
- CompletionResolver.java itself: UNCHANGED, zero diff (confirmed via
  `git diff --stat` before staging) -- its two call sites still call the
  2-arg onExhausted(target, reason), which is now the default and behaves
  identically. CompletionResolverTest (41 tests, 0 failures) proves this;
  its five ExhaustionSink lambdas needed a mechanical third parameter added
  to keep compiling against the new abstract method, no assertion changed.

Mutation proof, re-run against the new shape: forwardingTo's body edited to
call the 2-arg default instead of passing the hint through (the equivalent
of round 3's "delete the 3-arg override" now that there is only one method
to break) -- both new tests go red with the same assertions as round 3:

  ExhaustionSinkForwardingHazardTest...: expected: <gx> but was: <null>
  OpenCodeLauncherTest...ForwardingHop:  expected: <true> but was: <false>
  Tests run: 68, Failures: 2

Restored, re-ran: green (Tests run: 109, Failures: 0, including
CompletionResolverTest).

Compiler proof (not committed -- a scratch file outside the worktree,
compiled with the real ExhaustionSink.java on the classpath, then deleted):

  ExhaustionSink forwarder = (target, reason) -> System.out.println(target + reason);

  error: incompatible types: incompatible parameter types in lambda expression

A 2-arg lambda against this interface no longer compiles at all.

mvn clean install: Tests run: 1129, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-09-03 10:46:27 +07:00
Dai Ha 7840e9adf6 CB-201: make unmapped target notice truthful
CI / contract (pull_request) Successful in 1m13s
CI / build (pull_request) Successful in 1m15s
2026-09-03 10:44:18 +07:00
Dai Ha c3672f5472 fleetd #201/#227 unit 4: durable BACKEND_ERROR member outcome
CI / contract (push) Successful in 42s
CI / build (push) Successful in 1m47s
Adds MemberSession.State.BACKEND_ERROR with a nullable failureReason, surfaced in
rosterView, and SessionManager.onBackendError(target, reason). The CAS loop accepts
both sides of the completion race (BUSY and DONE); BACKEND_ERROR is terminal.

completeTurn now returns early when its CAS loses, so a stale DONE copy can no longer
release the pane or reset its context behind a member that just went BACKEND_ERROR.

Verified by the lead: 1130 tests green; mutating the completeTurn early return back to
the old fall-through turns losingCompletionDoesNotReleaseOrClearABackendErrorMember red
with 0 compile errors.
2026-09-03 10:43:58 +07:00
Dai Ha cf54aed451 CB-201 unit 2 review fix: threshold counts distinct targets, not raw events
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Successful in 1m46s
Two errors from the same target inside the window must never trip the
outage threshold on their own (a valid member report can legitimately
quote an "API Error:" line twice) — only two DIFFERENT targets on the
same credential do. Change evidenceCount() to targets.size() instead
of reasons.size(); reasons() still keeps every event, including
same-target repeats, so it can be longer than evidenceCount(). A real
outage still hits every target on the credential, so this loses no
true-positive coverage while cutting a real false-positive path.
2026-09-03 10:42:37 +07:00
Dai Ha bbf68f3e3c CB-201: cover losing completion CAS
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m37s
2026-09-03 10:35:36 +07:00
Dai Ha 826e0aeb2a CB-201 unit 2: credential-keyed backend outage policy
CI / build (pull_request) Successful in 1m11s
CI / contract (pull_request) Successful in 1m24s
Add BackendOutagePolicy: two classified backend errors on the same
credentialId within a 60s window mint one Incident and start a 60s
cool-off for that credential, one atomic ConcurrentHashMap.compute()
per credentialId so a concurrent second and third event can never
both cross the threshold. Errors during cool-off are ignored outright
(no extension, no incident); once cool-off elapses the next error
clears old evidence, requiring two fresh errors to rearm. This is a
new class, deliberately not BackendQuarantine (wrong store, wrong
1800s duration, misleading "exhausted" semantics for a 60s transient
fault). Knows nothing about panes, profiles, sessions, launchers, or
leads — takes events in, returns incidents out.
2026-09-03 10:30:40 +07:00
Dai Ha c935b181dd fleetd #234 round 3: make ExhaustionSink's forwarder a shared factory, not a rebuilt-per-caller shape
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m36s
Round 2's tests never reached Fleetd.java at all: both new tests declared
their OWN local copy of the forwarding shape instead of calling production's.
Mutating Fleetd.java's real forwarder back into the broken lambda left those
copies untouched, so the whole suite stayed green while production had
regressed to exactly the bug being fixed -- proven live by the reviewer.

Fix: extracted the forwarding shape into one named factory,
ExhaustionSink.forwardingTo(Supplier<ExhaustionSink> target), with the
"why a lambda here is wrong" explanation moved onto it (the one place the
shape is now written). Fleetd.java's forwarder collapses to one line:

    ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get);

Both new tests now call this same factory instead of rebuilding an anonymous
class inline, so they exercise the identical object production builds:
- ExhaustionSinkForwardingHazardTest: calls ExhaustionSink.forwardingTo
  directly and asserts the hint reaches the real sink through it.
- OpenCodeLauncherTest#theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop:
  same factory call, inside the full Fleetd-shaped construction order
  (forwarder built first, real sink pointed at via the AtomicReference
  afterward), driven through the real SessionManager.acquire() path.

Mutation proof, this time on production code only: deleted the factory's
3-arg override (falls back to the interface default, dropping the hint) --
both new tests go red with no test file touched:

  ExhaustionSinkForwardingHazardTest...: expected: <gx> but was: <null>
  OpenCodeLauncherTest...ForwardingHop:  expected: <true> but was: <false>
  Tests run: 68, Failures: 2

Restored, re-ran: green (Tests run: 68, Failures: 0). Confirmed Fleetd.java
carries no lambda ExhaustionSink anywhere (grep). ExhaustionSink.none() stays
a lambda on purpose -- both its overloads are true no-ops regardless of
arity, so there is no hint to drop.

mvn clean install: Tests run: 1129, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-09-03 10:30:37 +07:00
Dai Ha ee932fd85b CB-201: nudge leads about backend outages
CI / contract (pull_request) Successful in 1m7s
CI / build (pull_request) Successful in 1m55s
2026-09-03 10:28:51 +07:00
Dai Ha fe2e5ede34 CB-201: retain backend failure outcome
CI / build (pull_request) Successful in 1m11s
CI / contract (pull_request) Successful in 1m23s
2026-09-03 10:26:54 +07:00
Dai Ha c325054242 fleetd #234 round 2: fix the ExhaustionSink forwarding hop Fleetd.java actually uses
CI / contract (pull_request) Successful in 1m10s
CI / build (pull_request) Successful in 1m18s
The round-1 fix was dead on the real production path. Fleetd.java:177 builds
a forwarding sink (needed because the adapters are constructed before
`sessions` exists, breaking a genuine cycle) as a LAMBDA:

    ExhaustionSink forwardingExhaustionSink =
            (target, reason) -> exhaustionSinkRef.get().onExhausted(target, reason);

A lambda can only implement the interface's one abstract method (the 2-arg
overload), so it silently inherited the 3-arg overload's default body, which
drops the profile hint and calls back into the 2-arg method. OpenCodeLauncher
is constructed with this forwarder, so the hint it supplies (its own
already-known profile name) was thrown away before it ever reached the real
sink built later in Fleetd.main -- reproducing the exact silent no-op round 1
was sent to fix. The 1127 tests from round 1 all injected a sink directly
into OpenCodeLauncher and never went through this forwarding hop, so none of
them could see it.

Fix: forwardingExhaustionSink is now an anonymous class overriding both
overloads, each delegating to whatever exhaustionSinkRef currently holds.

Audited every other ExhaustionSink value in main/: the only other one is
ExhaustionSink.none() (a lambda), which is safe regardless of arity since
both its 2-arg body and the inherited 3-arg default are true no-ops.

New tests:
- ExhaustionSinkForwardingHazardTest: isolates the hazard at the interface
  level (a lambda forwarder drops the hint; an anonymous-class forwarder
  does not), independent of Fleetd.java's specific wiring.
- OpenCodeLauncherTest#theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop:
  replicates Fleetd.java's actual construction order (forwarder built and
  handed to the launcher first, real sink built and pointed at via the
  AtomicReference afterward) and drives the quarantine through it via the
  real SessionManager.acquire() path.

Both proven by mutation: temporarily rewriting each fixed forwarder back
into the pre-fix lambda makes its test fail with a real assertion message
(both matched exactly: "expected: <gx> but was: <null>" for the interface
proof, "expected: <true> but was: <false>" for the composed-wiring test);
restoring makes it pass again. No reverts were committed.

mvn clean install: Tests run: 1130, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-09-03 10:22:31 +07:00
Dai Ha 7662e2d0c8 fleetd #201/#227: refine backend outage work
CI / contract (push) Successful in 53s
CI / build (push) Successful in 1m49s
2026-09-03 10:21:44 +07:00
Dai Ha 4877992a70 fleetd #234: key the opencode model check on the resolved session id, and make the spawn-time quarantine actually happen
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Successful in 1m29s
Defect 1: OpenCodeSessionDiscovery.actualModelForDirectory queried
WHERE directory = ?, the same heuristic sessionIdForDirectory uses. Since a
default fleet_spawn (no worktree:) shares the lead's cwd with every other
worker and every past session ever run there, the model read-back could
silently compare against a DIFFERENT session's row. Renamed to
actualModelForSessionId(sessionId), keyed on the primary key id instead, and
made OpenCodeLauncher's SessionAwareHandle cache the resolved id once
non-null (AtomicReference) so a later sibling row in the same directory can
never flip which session's evidence is read. sessionIdForDirectory (#209) is
left directory-based on purpose, with a comment explaining why the heuristic
is unavoidable at that layer.

Defect 2: the ERROR log claimed "quarantining this profile's credential" but
Fleetd's ExhaustionSink lambda resolved target -> roster -> profile ->
credential, while OpenCodeLauncher's model-mismatch check fires from
agentSessionId() during SessionManager.acquire(), before the session is
registered in the roster -- the lookup found nothing and silently no-opped.
Added a default 3-arg ExhaustionSink.onExhausted(target, reason, profile)
overload (defaults to the 2-arg method, so CompletionResolver's two call
sites are unchanged); OpenCodeLauncher now passes its own already-known
profile name; Fleetd's sink became an anonymous class that tries the roster
first, falls back to the hint, and logs loudly at ERROR naming target/reason
when neither resolves, instead of silently no-oping.

Both fixes proven by mutation: reverting each independently makes its new
test fail with a real assertion message, restoring makes it pass again.

mvn clean install: Tests run: 1127, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-09-03 10:09:26 +07:00
ltms 1c051c4e47 fleetd #115: a profile that needs no token no longer warns, and no longer crashes
CI / contract (push) Successful in 1m2s
CI / build (push) Successful in 1m50s
tokenEnv no longer defaults to FLEETD_WORKER_TOKEN. A profile that names no
tokenEnv is stating it needs none, which is different from one naming a variable
that turns out to be unset. requiredSecretEnvVars now warns only for an explicitly
declared tokenEnv, so the permanent false alarm about a token no profile needs is
gone and the real warnings beside it stay trustworthy.

Decision recorded in a code comment: an explicit tokenEnv is checked for every
kind, opencode included. opencode can use its own provider credentials, but an
explicit tokenEnv declares a required host secret for its configured provider.

Second commit fixes a crash the first one introduced. Making tokenEnv nullable
changed what every reader of that value can receive, and ClaudeCodeLauncher:261
passed it straight to env.apply — System::getenv in production, which throws on a
null name. The live profile local-direct (kind: claude-code, baseUrl set, no
tokenEnv) would have crashed on spawn. It is weight: 0 today, so the failure would
have surfaced whenever someone re-enabled it. Now uses the superclass helper
resolveEnv, which already tolerates a null name and which OpenCodeLauncher was
already using.

Reader audit, all four: ClaudeCodeLauncher fixed; OpenCodeLauncher already safe via
resolveEnv; ConfigRef uses Objects.equals; MemberEnvAllowList drops null and blank
names in addIfPresent.

Verified by the lead before merge: reverting the resolveEnv fix makes the new test
fail with a NullPointerException from the null variable name, and restoring it
passes. Independent build: 1125 tests, 0 failures, 0 errors, 0 skipped.

Not verified: the ticket's acceptance criterion 4, a real boot on this host showing
no FLEETD_WORKER_TOKEN line while the WORKER_GITEA_TOKEN and AI_GATEWAY_TOKEN lines
are unchanged. A worker cannot restart the daemon it talks through. The lead checks
that at the next redeploy.
2026-09-03 05:02:27 +02:00
Dai Ha adb7a67880 CB-612: support tokenless Claude Code profiles
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Successful in 1m24s
2026-09-03 09:59:21 +07:00
Dai Ha 723fe494e9 CB-612: suppress unneeded token warnings
CI / contract (pull_request) Successful in 1m7s
CI / build (pull_request) Successful in 1m17s
2026-09-03 09:50:49 +07:00
25 changed files with 2388 additions and 121 deletions
+531
View File
@@ -0,0 +1,531 @@
# CB-201 and CB-227 refinement
Date: 2026-09-03
## Decision
#201 and #227 are one delivery program, but they are not one implementation unit.
#201 has a real seam: `CompletionResolver` can publish a typed backend-error event only after its
waiter resolution wins. #227 can consume that event without knowing any pane text. The classifier
must land before the final #227 wiring. However, the policy engine, roster state, and lead nudge can
be built in parallel with the classifier.
I propose five units. Units 1 to 4 own separate files and can run in parallel. Unit 5 owns all
composition files and lands after them. It also depends on the #234 defect 2 fix named in the task.
```mermaid
flowchart LR
U1["Unit 1: typed backend-error classification"] --> U5["Unit 5: wire policy, spawn gate, and fleet views"]
U2["Unit 2: credential outage policy"] --> U5
U3["Unit 3: lead outage nudge"] --> U5
U4["Unit 4: durable member outcome"] --> U5
D234["#234 defect 2: fail-loud target resolution"] --> U5
```
*Figure 1. Four file-disjoint foundations feed one composition unit.*
This split keeps `Fleetd.java` under one owner. It also keeps every other production file under one
unit in this plan.
## Evidence checked in the current branch
I read both issue pages in full. Each page reports zero comments.
| Evidence | What the code says now |
|---|---|
| `inject/CompletionResolver.java:229-237` | A turn below two seconds fails before normal scrape classification. A matching fast backend error is therefore only a generic failure today. |
| `inject/CompletionResolver.java:260-276` and `:332-367` | #211 already added raw-screen classification when `lastAssistantBlock` is empty. The “dead-code question” in #201 is stale on this branch. |
| `inject/CompletionResolver.java:288-317` | Exhaustion wins before the hard-coded `API Error:` match. A backend error then goes through generic `fail(...)`. |
| `inject/CompletionResolver.java:311-316` | The code admits that the pattern is a heuristic. A member report which quotes an API error may match it. |
| `inject/CompletionResolver.java:449-467` | Startup coverage exists only for `exhaustedPattern`. |
| `inject/ExhaustedPatternLookup.java:13-25` | The current lookup and explicit `none()` value are a good shape for the new classifier seam. |
| `Fleetd.java:196-207` | One `BackendQuarantine` is shared by placement and the exhaustion sink. Its cooldown comes from `quarantineCooldownSeconds`. |
| `Fleetd.java:322-363` | Pattern compilation, target-to-profile lookup, and the live `ExhaustionSink` are composed in `Fleetd.main`. The sink on this branch still ends in `.ifPresent(...)`. This plan assumes #234 replaces that silent path. |
| `placement/BackendQuarantine.java:60-87` | A repeated exhaustion restarts one long quarantine. The store is credential-keyed and uses an injected monotonic clock. |
| `member/CompositePeerLauncher.java:260-317` | Explicit and policy-selected spawns have separate gates. Both paths must learn about outage cool-off. |
| `member/CompositePeerLauncher.java:347-379` | Exhaustion refusal already checks a credential for explicit spawns and filters policy candidates. Its error text says “exhausted”. |
| `placement/PlacementContext.java:10-22` and `PlacementPolicyUtil.java:14-83` | Automatic placement has only one transient exclusion set named `quarantined`. Reusing it would make outage errors say “backend exhausted”. |
| `mcp/FleetMcp.java:913-1025` | `fleet_list` sets `free: 0` and adds `credentialId` plus `quarantinedForSeconds` when quarantine is active. |
| `session/MemberSession.java:51-59` | The roster has `DONE` and generic `FAILED`, but no backend-error state or stored reason. |
| `session/SessionManager.java:695-773` | A normal boundary moves `BUSY` to `DONE`. A failure moves any non-released session to `FAILED`. The async completion resolver can race the `DONE` update. |
| `session/SessionManager.java:648-687` | `rosterView` reports the session state, but it reports no terminal reason. |
| `msg/MessageService.java:922-940` | CB-588 already nudges for every terminal async ticket, including failures. Current code would report failed tickets, but it would not report one correlated outage. |
| `msg/ReplyPushLoop.java:20-48` | Replies, terminal tickets, and questions share one per-lead schedule. This prevents two push sources from injecting competing turns. |
| `msg/ReplyPushLoop.java:305-395` | Each push entry point resolves the owning lead through `PrimaryRegistry`. Missing ownership is logged and the durable or pending item remains the backstop. |
| `msg/ReplyPushLoop.java:496-547` | One tick builds one combined nudge. Pending items have separate reminder counts. |
| `health/FleetHealthMonitor.java:91-143` | Health is a slow periodic observer of members and message-layer facts. It does not receive completion classifications. |
| `health/FleetHealthMonitor.java:206-208` | `healthCoverage` means health enabled plus webhook configured. It does not describe lead-pane alerts. |
| `Fleetd.java:465-486` | Health stays `detection-only` without the webhook notification setting. |
I also read the related unit tests for `CompletionResolver`, `ReplyPushLoop`, `BackendQuarantine`,
`CompositePeerLauncher`, `PlacementPolicyUtil`, `SessionManager`, `MessageService`, and `FleetMcp`.
I did not inspect the in-progress #234 branch. I only used the two measured facts in the task. No
peer architect was named, so I did not exchange a design with one.
## Required behaviour
The policy should use these first values:
- Threshold: **2** classified backend errors.
- Window: **60 seconds**, measured from the first error to the second.
- Cool-off: **60 seconds**, starting when the threshold is reached.
- Correlation key: `credentialId`, never profile name and never error text.
- Incident rule: one active incident per credential. Errors during its cool-off do not extend it and
do not create more lead notices.
- Rearm rule: after cool-off ends, two fresh errors are needed for another incident.
Two errors are the smallest threshold which protects the honest one-turn failure. A 60-second window
fits the measured two-member outage. A 60-second cool-off blocks immediate repeat spawns without
turning a short backend fault into the default 1,800-second exhaustion quarantine.
A single classified error still fails its send and marks its member `backend_error`. It does not
cool a credential and does not send an outage notice. This is what “a single error changes nothing”
must mean at the credential level. It cannot mean that the failed member still looks successful.
```mermaid
sequenceDiagram
participant R1 as Resolver for member A
participant R2 as Resolver for member B
participant P as Outage policy
participant S as Spawn gate
participant N as Lead push loop
participant L as Lead pane
R1->>P: backend error for credential C
Note over P: Count 1, no cool-off
R2->>P: backend error for credential C within 60s
P->>P: Start one 60s incident
P->>S: Credential C is cooling off
P->>N: Queue one incident notice
N->>L: Inject when lead is idle, blocked, or done
L->>S: Request another spawn on credential C
S-->>L: Refuse and report remaining cool-off
```
*Figure 2. The second independent classification creates the fleet-level event.*
Against the 2026-09-01 case, the second failed member would start cool-off. `fleet_list` would show
zero free capacity and both members as `backend_error`. The push loop would inject one outage notice
even if the lead had not polled either ticket yet. The design reports the outage. It does not recover
uncommitted work from the members.
## Unit 1 — Typed backend-error classification
### Scope
Replace the direct hard-coded check inside `CompletionResolver` with a lookup and a sink. Keep the
public send result as a failed send. The typed internal event is the seam #227 consumes.
The lookup returns the pattern for a target. The sink receives the target, matched line, and full
failure reason. It fires only after `Rendezvous.resolveFailure(...)` wins for that exact captured
waiter. This copies the race rule already used by `ExhaustionSink`.
The classifier must run in all three current paths:
1. a normal non-empty assistant block;
2. the #211 raw scrape fallback;
3. a turn inside `MIN_TURN_NANOS`, before it becomes a generic too-fast failure.
In every path, the order stays: stale-baseline guard, exhaustion, backend error, then generic
failure or completion. A fast turn still fails when no configured pattern matches.
Keep `(?i)\bAPI Error\s*:` as a compatibility pattern for profiles without `errorPattern` until the
operator config is updated. Do not call this full coverage. Startup reporting in Unit 5 must name
profiles using this weaker legacy default.
### Files owned
- Add `fleetd/src/main/java/dev/ltms/fleet/inject/BackendErrorPatternLookup.java`.
- Add `fleetd/src/main/java/dev/ltms/fleet/inject/BackendErrorSink.java`.
- Change `fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java`.
- Change `fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java`.
No other unit may edit these files.
### Acceptance criteria
1. A target-specific error pattern matches a normal assistant block and resolves the send as failed.
2. The same match calls `BackendErrorSink` exactly once after the waiter resolution wins.
3. A late classification which loses to `fleet_reply` does not call the sink.
4. An exhausted line that also matches the generic error pattern stays `BACKEND_EXHAUSTED`. It calls
only `ExhaustionSink`.
5. A raw pane with leading Terminal User Interface (TUI) chrome and no assistant marker still uses
the #211 fallback and calls the backend-error sink.
6. A matching error inside the two-second floor is typed and sent to the sink. A non-matching fast
turn stays a generic failure.
7. An unchanged delivery baseline which contains old backend-error text is suppressed. It never
increments outage evidence.
8. A non-match keeps the existing completion result and text.
9. Constructors used by current callers keep compiling. They use the legacy default lookup and an
explicit inert sink until Unit 5 supplies the production objects.
10. Unit tests pass. The developer runs the focused test first, then `mvn clean install` from
`fleetd/`.
### Dependencies
None. Unit 1 can run with Units 2, 3, and 4.
Unit 5 depends on its new lookup, sink, and constructor.
### What to report back
- The exact classifier order in all three paths.
- The focused test command and result.
- The test which proves a losing waiter race does not publish an event.
- The test which proves a fast matching failure is typed.
- The final `mvn clean install` result.
- Any constructor kept only for transition and where Unit 5 replaces it.
## Unit 2 — Credential outage policy
### Scope
Build a small credential-keyed state machine. It accepts already-classified backend-error events.
It does not read pane text, profiles, sessions, or lead state.
Use an injected monotonic clock. A call records `credentialId`, target, and reason. It returns a new
incident only on the threshold crossing. The incident contains a stable event id, credential id,
the distinct affected targets, evidence count, window, and remaining cool-off.
This class owns both correlation and short cool-off. Keeping them together makes threshold crossing
and the cool-off deadline one atomic state change.
### Files owned
- Add `fleetd/src/main/java/dev/ltms/fleet/placement/BackendOutagePolicy.java`.
- Add `fleetd/src/test/java/dev/ltms/fleet/placement/BackendOutagePolicyTest.java`.
No other unit may edit these files.
### Acceptance criteria
1. One error creates no incident and no cool-off.
2. Two errors for one credential within 60 seconds create exactly one incident and a 60-second
cool-off.
3. Two errors more than 60 seconds apart do not create an incident.
4. The exact 60-second boundary has a pinned result. Use inclusive `<= 60s` so scheduler delay does
not discard evidence at the boundary.
5. Different credentials never share evidence.
6. Different profiles which supply the same credential id do share evidence. The policy itself only
sees the credential id.
7. More errors during active cool-off do not extend its deadline and do not return another incident.
8. After expiry, old evidence is cleared. Two fresh errors are needed to create the next incident.
9. Remaining seconds round up, matching `BackendQuarantine` reporting.
10. Concurrent second and third errors cannot return two incidents.
11. The focused tests and `mvn clean install` pass.
### Dependencies
None. Unit 2 can run with Units 1, 3, and 4.
Unit 5 depends on the policy API.
### What to report back
- The state transition table and locking method.
- The exact threshold, window, cool-off, and boundary rule.
- The test which proves one incident under concurrent calls.
- The focused test and `mvn clean install` results.
## Unit 3 — Lead outage nudge
### Scope
Add backend incidents as a fourth pending source in `ReplyPushLoop`. Do not create another scheduler
or call `AgentControl.send` from `Fleetd`. The existing combined per-lead schedule is the control
which prevents competing injected turns.
The entry point takes an incident id, affected worker targets, credential id, affected profile
names, and remaining cool-off. It resolves distinct owning leads through `PrimaryRegistry`.
Each `(incidentId, lead)` item is one-shot. It waits while the lead is not injectable. After one
successful `agents.send`, remove it. A send exception keeps it pending for a bounded retry. It never
uses the repeated reminder behaviour of an uncollected ticket.
Also add a fail-loud entry point for a classified target that Unit 5 cannot map to a credential. It
uses `PrimaryRegistry.nudgeTargetFor(target)` and says that correlation could not run. If no lead is
known, log at `WARN`, not `DEBUG`.
### Files owned
- Change `fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java`.
- Change `fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java`.
No other unit may edit these files.
### Acceptance criteria
1. One incident affecting two workers owned by one lead causes one successful pane injection.
2. Two affected workers owned by two leads cause one successful injection per affected lead. This
is one notice per event per lead, not one notice per member.
3. Repeating the same incident id is idempotent.
4. A busy or unknown lead is not injected. The item stays pending until the lead becomes injectable
or its attempt cap is reached.
5. After one successful injection, later ticks do not mention that incident again.
6. A failed `agents.send` is retried within the existing bound. A successful retry still gives only
one successful send.
7. A pending ticket and an outage incident for one lead appear in one combined nudge, not two
competing turns.
8. The text names the credential, profiles, affected workers, and remaining cool-off. It tells the
lead to run `fleet_list`.
9. An unmapped target produces a direct warning notice when a lead is known. If no lead is known,
the code logs a `WARN` naming the target and reason.
10. `stop()` clears incident state as it clears other push state.
11. Existing reply, ticket, and question tests stay green. The focused tests and
`mvn clean install` pass.
### Dependencies
None. The API uses plain values, not the Unit 2 incident class. This lets Unit 3 run in parallel.
Unit 5 adapts the Unit 2 incident into this entry point.
### What to report back
- The exact one-shot and retry rules.
- The test showing one combined nudge with a failed ticket.
- The test showing one successful send for two affected workers.
- The focused test and `mvn clean install` results.
## Unit 4 — Durable member backend-error outcome
### Scope
Make a classified backend failure remain visible after its ticket is collected or expires.
Add `BACKEND_ERROR` to `MemberSession.State`. Add a nullable failure detail to `MemberSession` and
render it as `failureReason` in `SessionManager.rosterView`. Add
`SessionManager.onBackendError(target, reason)`.
The transition must handle both completion orderings:
- `BUSY -> BACKEND_ERROR` when classification wins before the normal completion state update;
- `DONE -> BACKEND_ERROR` when the async resolver runs after `SessionManager.onTurnComplete`.
It must use a compare-and-set retry or another atomic update. `RELEASED` must never return to the
roster. A backend-error member is terminal and cannot accept another delivery.
### Files owned
- Change `fleetd/src/main/java/dev/ltms/fleet/session/MemberSession.java`.
- Change `fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java`.
- Change `fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java`.
No other unit may edit these files.
### Acceptance criteria
1. `onBackendError` moves a `BUSY` member to `BACKEND_ERROR` and stores the reason.
2. It also moves `DONE` to `BACKEND_ERROR`, covering the resolver race.
3. A later normal `onTurnComplete` cannot change `BACKEND_ERROR` back to `DONE`.
4. A released or unknown member is not recreated. The unknown case logs at `WARN` and returns an
explicit false result to its caller.
5. `onDelivered` refuses a `BACKEND_ERROR` member, just as it refuses generic `FAILED`.
6. `rosterView` reports `state: backend_error` and `failureReason` after the send ticket is gone.
7. Ordinary members do not gain a blank or invented `failureReason` field.
8. Existing constructors keep source compatibility for tests and adapters.
9. The focused tests and `mvn clean install` pass.
### Dependencies
None. Unit 4 can run with Units 1, 2, and 3.
Unit 5 calls the new session method from the production sink.
### What to report back
- The two race orderings and the tests for both.
- The exact roster JSON shape.
- The unknown-target result and log level.
- The focused test and `mvn clean install` results.
## Unit 5 — Production wiring, spawn gate, and fleet views
### Scope
Compose Units 1 to 4 in production. This is the only unit which edits `Fleetd.java`.
Add per-profile `errorPattern` config beside `exhaustedPattern`. Compile both once at startup. A
configured pattern wins over the legacy default. Report configured profiles and legacy-default
profiles separately at startup. A bad regex must stop startup with the profile and key in the
message.
Wire one production `BackendErrorSink` with this order:
1. mark the member `backend_error` with its reason;
2. resolve the profile and its current `effectiveCredentialId()` through the fail-loud #234 seam;
3. record the error in `BackendOutagePolicy`;
4. on a new incident, submit one event to `ReplyPushLoop`.
If target metadata cannot be resolved, do not end in `Optional.ifPresent`. Log an error and call the
Unit 3 unmapped-target notice. The failed send still reaches its ticket through CB-588.
Teach both spawn paths about a separate cool-off source. Exhaustion quarantine has priority when
both states are active. Automatic placement needs a distinct `coolingOff` set so its refusal does
not say “exhausted”.
Extend the MCP (Model Context Protocol) views:
- A cooling profile has `free: 0`, `credentialId`, and `coolingOffForSeconds` in `fleet_list`.
- It does not have `quarantinedForSeconds` unless exhaustion quarantine is also active.
- `fleet_profiles` has a separate `coolingOff` map, not an entry in `quarantined`.
- A direct spawn refusal says the credential is cooling off after repeated backend errors and gives
the remaining seconds.
Do not change `FleetHealthMonitor.coverage`. It still describes the periodic health webhook path.
Lead-pane outage delivery is a separate capability.
### Files owned
- `fleetd/src/main/java/dev/ltms/fleet/Fleetd.java`
- `fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java`
- `fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java`
- `fleetd/src/main/java/dev/ltms/fleet/member/CompositePeerLauncher.java`
- `fleetd/src/main/java/dev/ltms/fleet/placement/PlacementContext.java`
- `fleetd/src/main/java/dev/ltms/fleet/placement/PlacementPolicyUtil.java`
- `fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java`
- `fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java`
- `fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTest.java`
- `fleetd/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java`
- `fleetd/src/test/java/dev/ltms/fleet/placement/PlacementPolicyTest.java`
- `fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java`
- Add `fleetd/src/test/java/dev/ltms/fleet/BackendOutageFlowTest.java`.
- `fleetd/fleetd.example.yaml`
- `CLAUDE.md`
No earlier unit edits these files.
The lead, not a worker, must update `wiki/7-Use-Cases.md`, `wiki/9-Implementation.md`, and
`wiki/11-Features.md`. Project rules forbid workers from committing `wiki/`. The portable block in
`CLAUDE.md` and `wiki/7-Use-Cases.md` must remain byte-identical.
### Acceptance criteria
1. `errorPattern` binds per profile. Blank uses the legacy default and is reported as degraded
coverage. A config reload which changes it is reported as deferred because patterns are compiled
at startup.
2. A malformed `errorPattern` stops startup and names `profiles.<name>.errorPattern`.
3. The real production sink never silently drops an unknown target. A test captures its error log
and the fallback notice call.
4. One real backend-error classification through `CompletionResolver` marks only its member. It does
not cool the credential and does not send an outage notice.
5. Two real classifications for one credential within 60 seconds start one incident.
6. The integration test then calls the real explicit-profile spawn gate. It is refused before any
adapter spawn call, with “cooling off” and remaining seconds in the message.
7. The same test calls an automatic placement path. A cooling candidate is skipped. If every
candidate is cooling, the error names cool-off rather than exhaustion.
8. Profiles sharing the credential are all blocked. A profile on another credential stays usable.
9. `fleet_list` from the same fixture shows both members as `backend_error`, preserves each failure
reason, and reports `free: 0`, the credential, and `coolingOffForSeconds`.
10. `fleet_profiles` reports cool-off separately from quarantine.
11. The real `ReplyPushLoop` receives one incident and makes one successful lead-pane send. Existing
failed-ticket notice content may share that same combined send.
12. Exhaustion still wins when a line matches both patterns. A simultaneous exhaustion quarantine
also wins in spawn errors and fleet views.
13. After the 60-second cool-off, spawn is allowed again. A new incident needs two fresh errors.
14. `healthCoverage` has the same value before and after this change for the same health config.
15. `fleetd.example.yaml` explains `errorPattern`, the legacy fallback, 2/60/60 policy, and the
difference between cool-off and exhaustion quarantine.
16. `CLAUDE.md` tells leads how `fleet_profiles` and `fleet_list` report cool-off. The lead later
applies the matching wiki updates and runs the documented byte-sync check.
17. The developer records the new end-to-end test failing before implementation, then passing. The
focused suites and `mvn clean install` pass.
### Dependencies
Unit 5 starts only after Units 1 to 4 are merged or rebased into its branch. It also starts after the
#234 defect 2 fix lands, because both areas touch the same target-resolution control path.
### What to report back
- The exact commits used for Units 1 to 4 and #234.
- The startup coverage line with one configured and one legacy-default profile.
- The red test output before the implementation and its green result after.
- The explicit and automatic spawn refusal text.
- Sample `fleet_list` and `fleet_profiles` JSON for cool-off and exhaustion.
- The number and text of lead-pane sends in the real-path test.
- The focused test commands and final `mvn clean install` result.
- The exact `CLAUDE.md` change and the wiki edits the lead must apply.
## File ownership summary
| Area | Unit | Shared edit risk |
|---|---:|---|
| Completion classification | 1 | Only Unit 1 edits `CompletionResolver` and its test. |
| Correlation and cool-off state | 2 | New files only. |
| Lead push scheduling | 3 | Only Unit 3 edits `ReplyPushLoop` and its test. |
| Member terminal state | 4 | Only Unit 4 edits `MemberSession`, `SessionManager`, and their test. |
| Main composition, config, placement, MCP views, shipped prompt | 5 | Only Unit 5 edits `Fleetd`, `FleetConfig`, `CompositePeerLauncher`, placement context, `FleetMcp`, and `CLAUDE.md`. |
| Wiki propagation | Lead after Unit 5 | Workers do not commit the wiki submodule. |
## What I would not build
1. **Do not reuse `BackendQuarantine` for outages.** Its repeat call restarts a long credential
quarantine. Its fields and errors say “exhausted”. That is wrong for a short outage.
2. **Do not merge the exhaustion and generic error patterns.** Exhaustion must win because it has a
different policy and duration.
3. **Do not group by error string.** One outage can produce different text. The shared operational
limit is the credential.
4. **Do not mark a profile unusable until config changes.** The current classifier cannot safely
tell a permanent malformed request from a transient service fault. A permanent state would need
a stronger error taxonomy first.
5. **Do not quarantine on the first generic backend error.** That would turn one bad request or one
false pattern match into a fleet-wide capacity loss.
6. **Do not add this to `FleetHealthMonitor`.** The monitor samples slow member health. The exact
backend event already exists at completion resolution, and moving it to polling would lose type
and time.
7. **Do not add another direct lead injector.** `ReplyPushLoop` already owns status gating,
per-lead coalescing, retry bounds, and heartbeat stand-down.
8. **Do not change `healthCoverage` to `full`.** That field still means a webhook notification sink
exists for periodic health. A backend outage nudge does not make every health event visible.
9. **Do not persist incident history across daemon restart in this work.** Existing exhaustion
quarantine is also in memory. A 60-second state does not justify a new durable store.
10. **Do not build work recovery.** The PR-body survival story proves why checkpoint-first work is
useful, but these tickets are about detection, capacity, and signalling.
11. **Do not remove the legacy `API Error:` fallback in the first release.** Doing so would turn an
unedited config back into a false successful completion. Report it as degraded coverage instead.
12. **Do not reorder or add the old `visibleTurn` fallback from #201.** #211 already implemented the
narrow raw-scrape fallback at `CompletionResolver.classifyRawScrapeFallback`.
## Riskiest assumption and cheapest experiment
The riskiest assumption is that a configured error regex means “the backend failed this turn”. The
current code and test already show the counterexample: a worker may quote `API Error:` while writing
a valid report. Two such false matches on one credential would now remove capacity for 60 seconds.
The cheapest experiment is a replay corpus before Unit 5 ships:
1. Save the full pane text from the measured 2026-09-01 outage.
2. Produce one safe failure per backend with a disposable invalid endpoint or request.
3. Save one valid member report which quotes each error line.
4. Replay all samples through the real `CompletionResolver` test fixture.
5. Require outage samples to match and quoted-report samples not to match after assistant-block
extraction and baseline checks.
This costs no outage deployment and no real sleep. If quoted reports still match, narrow the profile
patterns before enabling correlation. Do not raise the threshold to hide a bad classifier.
## Sequencing with three developers
First wave:
1. Developer A: Unit 1, typed classification.
2. Developer B: Unit 2, credential outage policy.
3. Developer C: Unit 3, lead outage nudge.
As soon as one slot is free, start Unit 4. It is file-disjoint from every first-wave unit. Merge and
review Units 1 to 4 independently.
Start Unit 5 only after all four foundations and #234 are available. Unit 5 is the only high-conflict
integration branch, so no other active unit should touch its file list.
## Checks performed for this refinement
- Read issue #201 and issue #227 through their Gitea pages. Both showed zero comments.
- Read the source and tests named in the evidence section.
- Ran `git status --short --branch`; the branch was clean before this document was added.
- Ran `git log --oneline -12` to identify the branch base.
- I did not run Maven because this change adds only a design document.
- Rendered both Mermaid blocks with `npx @mermaid-js/mermaid-cli`; both commands succeeded.
+53 -16
View File
@@ -174,7 +174,13 @@ public final class Fleetd {
// 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 = (target, reason) -> exhaustionSinkRef.get().onExhausted(target, reason);
// 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).
ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
// 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.
@@ -345,22 +351,51 @@ public final class Fleetd {
// model-mismatch check — it needs nothing profile-specific from the caller beyond `target`
// (a herdr terminal id) and `reason`, so reusing it here is exactly "the existing
// ExhaustionSink path", not a new mechanism.
ExhaustionSink exhaustionSink = (target, reason) -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(MemberSession::profile)
.map(profileName -> config.get().profiles().get(profileName))
.ifPresent(profile -> {
String credentialId = profile.effectiveCredentialId();
quarantine.quarantine(credentialId);
log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId,
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
});
//
// fleetd #234: that check fires from SessionAwareHandle.agentSessionId(), which runs during
// SessionManager.acquire() BEFORE this session is registered in sessions.roster() — so the
// roster-only lookup below used to find nothing, .ifPresent silently no-op'd, and the
// ERROR the check had just logged ("quarantining this profile's credential") was a lie:
// nothing was quarantined, and nothing said so. Two changes: (1) OpenCodeLauncher now
// passes its OWN profile name via ExhaustionSink's 3-arg overload — it already has the
// FleetConfig.Profile in hand and does not need the roster at all — used here as a
// fallback whenever the roster lookup misses; (2) if a profile still cannot be resolved
// (neither the roster nor the hint names a configured one), this logs loudly at ERROR
// instead of silently doing nothing — a control that cannot act must say so.
// fleetd #234, round 4: the 3-arg overload is now ExhaustionSink's single abstract method,
// so this is safely a lambda — there is no separate 2-arg overload left for it to bind to
// instead and silently drop profileHint (that was rounds 1-3's whole hazard).
ExhaustionSink exhaustionSink = (target, reason, profileHint) -> {
String profileName = sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(MemberSession::profile)
.orElse(profileHint);
FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName);
if (profile == null) {
log.error("quarantine requested for target '{}' ({}) but no profile could be "
+ "resolved — the target is not (yet) in the roster, and {} — "
+ "credential NOT quarantined (fleetd #234)",
target, reason,
profileHint == null ? "no profile hint was given"
: "the hinted profile '" + profileHint + "' is not configured");
return;
}
String credentialId = profile.effectiveCredentialId();
quarantine.quarantine(credentialId);
log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId,
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
};
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
// now that `sessions` exists to resolve target -> session -> profile.
exhaustionSinkRef.set(exhaustionSink);
AgentControl agents = router.memberAgents();
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink);
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink,
target -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> new CompletionResolver.WorktreeBranch(session.worktree(), session.branch()))
.orElse(null));
// 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();
@@ -795,7 +830,7 @@ public final class Fleetd {
/**
* CB-594: which env vars the loaded config actually needs, and why — every non-{@code
* subscription} profile's {@code tokenEnv} (a subscription profile never reads one, see
* subscription} profile's explicitly configured {@code tokenEnv} (a subscription profile never reads one, see
* {@link FleetConfig.Profile#isSubscription()}), plus every profile's {@code gitTokenEnv}
* where set (opt-in), plus a configured {@code broker.uriEnv} (CB-151). Derived from the
* config, not hard-coded, so a new profile is covered for free. A var required by more than one
@@ -812,7 +847,9 @@ public final class Fleetd {
static Map<String, List<String>> requiredSecretEnvVars(FleetConfig cfg) {
Map<String, List<String>> requiredBy = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (!profile.isSubscription()) {
// All kinds are checked when tokenEnv is explicit. OpenCode may use provider credentials,
// but an explicit tokenEnv still declares a required host secret for its configured provider.
if (!profile.isSubscription() && profile.hasTokenEnv()) {
requiredBy.computeIfAbsent(profile.tokenEnv(), _ -> new ArrayList<>())
.add("profile '" + name + "' tokenEnv");
}
@@ -838,7 +875,7 @@ public final class Fleetd {
* <p>A missing entry only warns — it must never refuse to start. A daemon that boots and says
* what is wrong is strictly more useful than one that will not boot at all.
*/
private static void reportRequiredSecrets(FleetConfig cfg) {
static void reportRequiredSecrets(FleetConfig cfg) {
Map<String, List<String>> requiredBy = requiredSecretEnvVars(cfg);
if (requiredBy.isEmpty()) {
log.info("startup secrets: no profile references a token env var — nothing to check");
@@ -240,7 +240,8 @@ public record FleetConfig(
* @param configDir {@code CLAUDE_CONFIG_DIR} so the worker inherits the profile's
* skills/MCP/hooks (may be {@code null})
* @param tokenEnv name of the host env var holding the worker's auth token; its value
* is injected as {@code ANTHROPIC_AUTH_TOKEN} (never stored in config)
* is injected as {@code ANTHROPIC_AUTH_TOKEN} (never stored in config).
* {@code null}/blank means this profile needs no token.
* @param argv launch command; defaults to {@code ["claude"]}
* @param placement where a worker lands: {@code "tab"} (default — its own tab in the
* worker space) or {@code "pane"} (legacy — split the focused tab)
@@ -389,7 +390,7 @@ public record FleetConfig(
? (KIND_CLAUDE_CODE.equals(k) ? List.of("claude") : List.of(k))
: List.copyOf(argv);
kind = k;
tokenEnv = (tokenEnv == null || tokenEnv.isBlank()) ? "FLEETD_WORKER_TOKEN" : tokenEnv;
tokenEnv = (tokenEnv == null || tokenEnv.isBlank()) ? null : tokenEnv;
placement = (placement == null || placement.isBlank()) ? "tab" : placement.toLowerCase();
workspace = (workspace == null || workspace.isBlank()) ? "fleet" : workspace;
// CB-557: no per-profile default any more. A label is generated from the member's ROLE
@@ -604,6 +605,11 @@ public record FleetConfig(
return gitTokenEnv != null && !gitTokenEnv.isBlank();
}
/** True when this profile explicitly names a host token environment variable. */
public boolean hasTokenEnv() {
return tokenEnv != null && !tokenEnv.isBlank();
}
/**
* True when this profile's {@code env:} block names a worker-side Anthropic binding variable.
* Those two keys are the adapter's, never the operator's: on a claude-code worker
@@ -14,6 +14,7 @@ import java.util.TreeSet;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.LongSupplier;
import java.util.function.Function;
import java.util.regex.Pattern;
/**
@@ -85,6 +86,17 @@ public final class CompletionResolver implements TurnListener {
private static final String CLIPPED_PANE_TAIL_MARKER =
"[Pane tail clipped: member did not call fleet_reply.]";
/** A pane echo must be this large before it can replace a completion report. */
static final int ECHO_MIN_CHARS = 400;
/** Normalised TUI chrome may add this many characters to an otherwise echoed brief. */
static final int MAX_ECHO_EXCESS_CHARS = 160;
/** The explicit result returned instead of a lead's echoed injected brief. */
public static final String NO_REPORT_PREFIX = "[no report — the member ended its turn without fleet_reply, "
+ "and the pane still shows the injected brief. Nothing was produced on the pane. Check the "
+ "member's worktree and branch for committed work before re-delegating.";
private final AgentControl agents;
private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
@@ -92,6 +104,11 @@ public final class CompletionResolver implements TurnListener {
private final BackendErrorPatternLookup backendErrorPatterns;
private final BackendErrorSink backendErrorSink;
private final LongSupplier nowNanos;
private final Function<String, WorktreeBranch> worktreeBranches;
/** Known member location, used only to guide a lead after an echoed brief. */
public record WorktreeBranch(String worktree, String branch) {
}
/**
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
@@ -108,7 +125,8 @@ public final class CompletionResolver implements TurnListener {
* the other half of the {@link #MIN_TURN_NANOS} floor check, compared against a fresh reading at
* resolution time.
*/
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos,
String injectedText) {
/**
* Convenience for tests exercising scrape/suppression logic that don't care about turn
@@ -116,7 +134,11 @@ public final class CompletionResolver implements TurnListener {
* Not used by production code — {@link #captureBaseline} always records a real reading.
*/
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
this(waiter, baseline, Long.MIN_VALUE / 2);
this(waiter, baseline, Long.MIN_VALUE / 2, null);
}
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
this(waiter, baseline, deliveredAtNanos, null);
}
}
@@ -139,11 +161,18 @@ public final class CompletionResolver implements TurnListener {
* (fleetd#201 Unit 5).
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink) {
ExhaustionSink exhaustionSink) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink,
BackendErrorPatternLookup.legacy(), BackendErrorSink.none(), System::nanoTime);
}
/** Production constructor with a lookup for the member worktree and branch. */
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, Function<String, WorktreeBranch> worktreeBranches) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink,
BackendErrorPatternLookup.legacy(), BackendErrorSink.none(), System::nanoTime, worktreeBranches);
}
/**
* Transition constructor (fleetd#201 Unit 1): same legacy backend-error defaults as the 4-arg
* constructor above, but with the injectable clock. Kept so existing fleetd#164 timing tests
@@ -172,7 +201,7 @@ public final class CompletionResolver implements TurnListener {
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, backendErrorPatterns, backendErrorSink,
System::nanoTime);
System::nanoTime, _ -> null);
}
/**
@@ -185,8 +214,17 @@ public final class CompletionResolver implements TurnListener {
* package.
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink, LongSupplier nowNanos) {
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink, LongSupplier nowNanos) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, backendErrorPatterns, backendErrorSink,
nowNanos, _ -> null);
}
/** Full constructor with injectable clock and member location lookup. */
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink, LongSupplier nowNanos,
Function<String, WorktreeBranch> worktreeBranches) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
@@ -194,6 +232,7 @@ public final class CompletionResolver implements TurnListener {
this.backendErrorPatterns = Objects.requireNonNull(backendErrorPatterns, "backendErrorPatterns");
this.backendErrorSink = Objects.requireNonNull(backendErrorSink, "backendErrorSink");
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
this.worktreeBranches = Objects.requireNonNull(worktreeBranches, "worktreeBranches");
}
@Override
@@ -223,7 +262,7 @@ public final class CompletionResolver implements TurnListener {
baseline = null; // fail open: no baseline ⇒ no suppression
log.debug("delivery baseline for {} failed: {}", target, e.getMessage());
}
inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong()));
inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong(), token.injectedText()));
}
/** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */
@@ -369,6 +408,9 @@ public final class CompletionResolver implements TurnListener {
return;
}
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (echoesInjectedBrief(tail, turn.injectedText())) {
completion = noReportMessage(target) + (clipped ? "\n" + CLIPPED_PANE_TAIL_MARKER : "");
}
if (rendezvous.resolveCompletion(waiter, completion)) {
inFlight.remove(target, turn);
if (clipped) {
@@ -381,6 +423,40 @@ public final class CompletionResolver implements TurnListener {
}
}
/**
* A full echoed brief is at least 400 normalised characters. A scrape that contains the brief may
* add no more than 160 normalised characters of TUI chrome. This accepts harmless status text, but
* preserves a real report that restates the full brief before adding substantive content.
*/
static boolean echoesInjectedBrief(String scrape, String injectedText) {
String normalScrape = normalize(scrape);
String normalInjected = normalize(injectedText);
if (normalScrape.length() < ECHO_MIN_CHARS || normalInjected.length() < ECHO_MIN_CHARS) {
return false;
}
if (normalInjected.contains(normalScrape)) {
return true;
}
return normalScrape.contains(normalInjected)
&& normalScrape.length() - normalInjected.length() <= MAX_ECHO_EXCESS_CHARS;
}
private static String normalize(String text) {
return text == null ? "" : text.toLowerCase().replaceAll("[^a-z0-9]+", "");
}
private String noReportMessage(String target) {
WorktreeBranch location = worktreeBranches.apply(target);
if (location == null || (location.worktree() == null && location.branch() == null)) {
return NO_REPORT_PREFIX + "]";
}
String locationText = location.worktree() == null ? "" : " worktree=" + location.worktree();
if (location.branch() != null) {
locationText += " branch=" + location.branch();
}
return NO_REPORT_PREFIX + locationText + "]";
}
/**
* fleetd#211: the raw-scrape fallback classification, run only when {@link #lastAssistantBlock}
* found nothing usable (see the call site in {@link #resolve}). Mirrors the two classifications
@@ -1,5 +1,7 @@
package dev.ltms.fleet.inject;
import java.util.function.Supplier;
/**
* Notified when {@link CompletionResolver} actually delivers a {@code BACKEND_EXHAUSTED}
* classification to a waiting send (CB-578 stage B) — never on a race that lost (see
@@ -10,22 +12,81 @@ package dev.ltms.fleet.inject;
* profiles or credentials, so mapping {@code target} to whatever should be quarantined is entirely
* the sink's job — see {@code Fleetd.main}'s wiring, which resolves target → session → profile →
* {@code effectiveCredentialId()} and calls {@code BackendQuarantine.quarantine} on it.
*
* <p><strong>The 3-arg overload is the single abstract method — fleetd #234, round 4.</strong> Two
* earlier rounds each shipped a caller that silently dropped the profile hint (see {@link
* #onExhausted(String, String, String)}): a lambda written against this interface can only ever
* implement whichever overload is abstract, and while the 2-arg form held that position, EVERY
* lambda site — a call-site forwarder in {@code Fleetd.java}, a hand-built test double — bound to
* it and silently inherited the profile-dropping default, whether or not its author remembered the
* hazard. Making the 3-arg form abstract instead removes the shape entirely: a lambda declared
* against this interface today is <em>forced</em> by the compiler to take {@code (target, reason,
* profile)}, so there is no overload left for it to bind to that can drop the hint. This is a type
* change, not a test — it holds even for a caller that has never heard of fleetd #234.
*/
@FunctionalInterface
public interface ExhaustionSink {
/**
* @param target the herdr terminal id whose turn was classified {@code BACKEND_EXHAUSTED}
* @param reason the matched-line reason carried by the classification
* @param target the herdr terminal id whose turn was classified {@code BACKEND_EXHAUSTED}
* @param reason the matched-line reason carried by the classification
* @param profile the profile the caller already knows should be quarantined, or {@code null}
* when the caller has no better answer than {@code target} alone (fleetd #234):
* a caller whose {@code target} is not yet resolvable through whatever roster
* the sink's implementation consults — {@link
* dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch check fires from
* {@code SessionAwareHandle.agentSessionId()}, which runs during {@code
* SessionManager.acquire()} <em>before</em> that session is registered, so a
* target -> session -> profile lookup finds nothing at that point. That launcher
* already has its own {@code FleetConfig.Profile} in hand and does not need the
* roster to know which profile to quarantine, so it supplies this directly
* instead of leaving the sink to guess.
*/
void onExhausted(String target, String reason);
void onExhausted(String target, String reason, String profile);
/**
* Convenience for a caller with no profile to offer — every existing call site that predates
* the hint (fleetd #234): {@link dev.ltms.fleet.inject.CompletionResolver}'s two call sites
* always call with a {@code target} that IS live in the roster at the time of the call, so they
* need no hint and keep working exactly as before, unchanged by this default.
*
* @param target as {@link #onExhausted(String, String, String)}
* @param reason as {@link #onExhausted(String, String, String)}
*/
default void onExhausted(String target, String reason) {
onExhausted(target, reason, null);
}
/**
* Inert sink — nothing happens on exhaustion. The explicit stand-in a caller (or a test not
* exercising this feature) passes instead of a defaulting overload, exactly like
* {@link ExhaustedPatternLookup#none()}.
* {@link ExhaustedPatternLookup#none()}. Safe as a lambda: the 3-arg form is now the interface's
* single abstract method, so a lambda here has no other overload to silently bind to instead —
* it simply does nothing with all three arguments.
*/
static ExhaustionSink none() {
return (target, reason) -> { };
return (target, reason, profile) -> { };
}
/**
* A sink that forwards to whatever {@code target} currently supplies (fleetd #234). Exists to
* break a genuine construction-order cycle: {@code Fleetd.main} builds its adapters (including
* {@link dev.ltms.fleet.member.OpenCodeLauncher}) before {@code sessions} exists, so it cannot
* hand them the real sink yet — it hands them a forwarder pointed at an {@code
* AtomicReference<ExhaustionSink>} that starts at {@link #none()} and gets {@code .set()} to the
* real sink once {@code sessions} is built. {@code target} is evaluated on every call, never
* cached, so the forwarder keeps working after the reference is repointed.
*
* <p>Now safe as a one-line lambda (round 4): forwarding the single 3-arg abstract method
* forwards everything a caller can supply — there is no separate 2-arg overload left for a
* forwarder to bind to instead and silently lose the hint. Kept as a named factory rather than
* written inline at each call site anyway, so a test can call the exact object {@code
* Fleetd.java} builds instead of asserting a rebuilt copy of its shape (round 3's lesson).
*
* @param target supplies the sink to forward to, evaluated fresh on every call
* @return a sink whose call delegates to {@code target.get()}
*/
static ExhaustionSink forwardingTo(Supplier<ExhaustionSink> target) {
return (t, r, p) -> target.get().onExhausted(t, r, p);
}
}
@@ -10,6 +10,7 @@ import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadMessage;
@@ -543,8 +544,9 @@ public final class FleetMcp {
case REPLIED -> text(r.text());
// The worker's turn finished but it never called fleet_reply — hand back the scraped
// transcript tail, flagged so the primary knows it isn't a structured reply.
case COMPLETED_UNREPLIED -> text(
"[worker finished without a structured fleet_reply — transcript tail follows]\n" + r.text());
case COMPLETED_UNREPLIED -> text(r.text().startsWith(CompletionResolver.NO_REPORT_PREFIX)
? r.text()
: "[worker finished without a structured fleet_reply — transcript tail follows]\n" + r.text());
// The worker ran the turn then wedged (CB-109) — surface the error context.
case WORKER_FAILED -> text("[worker failed — turn ended in an unrecoverable state]\n" + r.text());
// The backend refused on a subscription usage limit (CB-578 stage A) — the worker's
@@ -659,6 +661,7 @@ public final class FleetMcp {
}
return switch (v.phase()) {
case DONE -> text(v.replySource() != null && v.replySource().equals("transcript")
&& !v.reply().startsWith(CompletionResolver.NO_REPORT_PREFIX)
? "[done — worker finished without a structured fleet_reply; transcript tail follows]\n" + v.reply()
: v.reply());
case PENDING -> text("[pending — " + v.detail() + "]");
@@ -258,7 +258,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
workerEnv.remove("ANTHROPIC_AUTH_TOKEN");
} else {
workerEnv.put("ANTHROPIC_BASE_URL", baseUrl);
putIfPresent(workerEnv, "ANTHROPIC_AUTH_TOKEN", env.apply(cfg.tokenEnv()));
putIfPresent(workerEnv, "ANTHROPIC_AUTH_TOKEN", resolveEnv(cfg.tokenEnv()));
}
putIfPresent(workerEnv, "ANTHROPIC_MODEL", cfg.model());
putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir());
@@ -22,6 +22,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.BooleanSupplier;
import java.util.function.Function;
import java.util.function.LongSupplier;
@@ -706,6 +707,17 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
private final ExhaustionSink exhaustionSink;
/** CAS'd true the first (and only) time a model mismatch is reported for this handle. */
private final AtomicBoolean modelMismatchReported = new AtomicBoolean();
/**
* The session id, once {@link OpenCodeSessionDiscovery#sessionIdForDirectory} first
* resolves a non-null answer for this handle (fleetd #234). Sticky on purpose: {@code
* directory} is a shared-cwd heuristic (see {@link OpenCodeSessionDiscovery}'s class
* javadoc) that can start returning a DIFFERENT row once another session shares the same
* directory and writes a newer one — re-deriving it on every call would let this handle's
* identity silently drift to a sibling's session. Once resolved, this IS the answer, and
* {@link #checkModelMatch} reads only the row this id names, never "whatever is newest in
* the directory right now."
*/
private final AtomicReference<String> resolvedSessionId = new AtomicReference<>();
SessionAwareHandle(PeerHandle delegate, OpenCodeSessionDiscovery discovery, String cwd,
FleetConfig.Profile cfg,
@@ -759,17 +771,27 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
}
return null;
}
// fleetd #234: once resolved, stay resolved. Re-deriving from `directory` on every call
// would let this handle's identity drift to a sibling session that later shares the
// same cwd and writes a newer row — see resolvedSessionId's javadoc.
String cached = resolvedSessionId.get();
if (cached != null) {
return cached;
}
// Lazy + retried, never a spawn-time blocker: opencode writes the session record only
// when the session is first persisted, so null here is the correct interim answer and
// the caller re-calls later (each call re-scans, picking up a record that has since
// appeared).
String id = discovery.sessionIdForDirectory(cwd);
if (id != null) {
resolvedSessionId.compareAndSet(null, id);
}
// fleetd #175: check on the SAME tick — while the caller (SessionManager's late-resolve
// step) is still re-polling because the id is unknown, the row this id came from (once
// it exists) is exactly the row that also carries the actual model. Once id resolves,
// the caller stops calling agentSessionId() for this session, so this is naturally a
// once-only check that happens right when the row first appears.
checkModelMatch();
checkModelMatch(id);
return id;
}
@@ -777,16 +799,25 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* Verify the live opencode session is running the model {@link #cfg} requested (fleetd
* #175) and, on a real mismatch, log an ERROR and quarantine through {@link
* #exhaustionSink}. A no-op when there is nothing to compare against — no model configured,
* already reported once for this handle, or the actual model is still UNKNOWN (no row yet,
* unreadable database, or unparseable evidence). UNKNOWN must never be treated as a
* mismatch: that is the single most important safety rule here — a false positive would
* already reported once for this handle, {@code sessionId} itself is not resolved yet
* (fleetd #234: absent evidence, not a mismatch), or the actual model is still UNKNOWN (no
* row yet, unreadable database, or unparseable evidence). UNKNOWN must never be treated as
* a mismatch: that is the single most important safety rule here — a false positive would
* quarantine a perfectly working profile's credential.
*
* @param sessionId the id {@link #agentSessionId()} just resolved (or had cached) for THIS
* handle — the model is read back for this exact session (fleetd #234's
* {@link OpenCodeSessionDiscovery#actualModelForSessionId}), never
* re-derived from {@code directory}
*/
private void checkModelMatch() {
private void checkModelMatch(String sessionId) {
if (modelMismatchReported.get() || cfg.model() == null || cfg.model().isBlank()) {
return;
}
OpenCodeSessionDiscovery.ActualModel actual = discovery.actualModelForDirectory(cwd);
if (sessionId == null || sessionId.isBlank()) {
return; // id not resolved yet — UNKNOWN, never a mismatch (fleetd #175's rule)
}
OpenCodeSessionDiscovery.ActualModel actual = discovery.actualModelForSessionId(sessionId);
if (actual == null) {
return; // UNKNOWN evidence — never a mismatch
}
@@ -817,10 +848,16 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
+ "falls back to a default model, which may be a PAID credential "
+ "(fleetd #175); quarantining this profile's credential",
cfg.profile(), cfg.model(), actualDisplay);
// fleetd #234: pass our OWN profile name too. This check fires from agentSessionId(),
// called during SessionManager.acquire() BEFORE this session is registered in
// sessions.roster() — a roster-only sink (Fleetd's target -> session -> profile lookup)
// finds nothing at this point and silently no-ops (defect 2). We already know exactly
// which profile to quarantine without the roster; the sink is passed it explicitly.
exhaustionSink.onExhausted(delegate.terminalId(),
"opencode model mismatch: profile '" + cfg.profile() + "' requested '"
+ cfg.model() + "' but the live session is running '" + actualDisplay
+ "' (fleetd #175)");
+ "' (fleetd #175)",
cfg.profile());
}
@Override
@@ -32,11 +32,19 @@ import java.util.concurrent.atomic.AtomicBoolean;
* here, so a layout change, or a switch to the HTTP server, changes exactly one class and nothing
* in {@link OpenCodeLauncher}.
*
* <p>The determinism that makes this useful is structural, not a guess: every fleetd worker runs
* in its own unique git worktree, so the row's {@code directory} (its project root) equals the
* worker's cwd identifies <em>its</em> session unambiguously. We match on {@code directory} rather
* than diffing {@code opencode session list} before/after — that races under concurrent spawns, and
* the CLI listing does not even show the directory.
* <p><strong>The {@code directory} match is a heuristic, not an identity — fleetd #234.</strong> A
* worktree is opt-in: {@code fleet_spawn} only provisions one when the caller passes {@code
* worktree:}; the default spawn inherits the lead's own cwd, which every other worker spawned the
* same way (and every past session ever run there) shares. {@code directory} therefore does
* <em>not</em> identify a session unambiguously in general — only in the special case of a fresh,
* unique worktree does "most recently updated row for this directory" reliably mean "this worker's
* own row." {@link #sessionIdForDirectory} still has to use this heuristic (the id has to come from
* somewhere, and nothing else is available at this layer — see that method's javadoc), but a caller
* that already holds a resolved id must never re-derive evidence about that same session via
* {@code directory} again; see {@link #actualModelForSessionId}, which looks up by {@code id}
* instead for exactly this reason. We match on {@code directory} rather than diffing
* {@code opencode session list} before/after — that races under concurrent spawns, and the CLI
* listing does not even show the directory.
*
* <p>All reads are best-effort and never throw: a missing or unreadable database, a query that
* fails, or a directory with no row yet all yield {@code null}, and the caller (the session
@@ -89,8 +97,16 @@ final class OpenCodeSessionDiscovery {
/**
* The opencode session id whose row references {@code directory} (the worker's cwd), or
* {@code null} when no row matches yet. When several rows share the directory — e.g. repeated
* spawns into the same worktree — the row with the highest {@code time_updated} wins: it is
* the session the pane most likely corresponds to.
* spawns into the same worktree, OR several workers sharing one cwd because none of them was
* given a worktree (fleetd #234) — the row with the highest {@code time_updated} wins: it is
* the session the pane most likely corresponds to. That "most likely" is a real caveat, not a
* formality: when the directory is shared, this can and does pick another session's row (see
* the class javadoc). This is the one place fleetd resolves an opencode session id at all —
* nothing else is available at this layer to disambiguate further (no {@code opencode session
* list} entry names the directory, and diffing before/after races under concurrent spawns) — so
* the heuristic stays here unchanged. What must never happen is a SECOND, independent piece of
* evidence about the same session being re-derived via {@code directory} once an id has already
* come out of this method; see {@link #actualModelForSessionId}.
*
* <p>Never throws: a missing {@code opencode.db}, a locked/unreadable database, a query
* failure, or a directory that has not been persisted yet all resolve to {@code null} rather
@@ -134,24 +150,36 @@ final class OpenCodeSessionDiscovery {
}
/**
* The model opencode actually ran the {@code directory}'s most-recent session on (fleetd
* #175), read from the same row {@link #sessionIdForDirectory} matches — but via its own
* query and its own connection, deliberately kept independent so a database whose schema
* predates the {@code model} column (or any other read failure on this column alone) can
* never take {@link #sessionIdForDirectory}'s id resolution down with it. That would be a
* regression of the id-resolution feature #209 shipped; this method degrades on its own.
* The model opencode actually ran the session {@code sessionId} on (fleetd #175/#234), read by
* primary-key lookup — the ONE row that id names, and no other. Deliberately keyed on
* {@code id} rather than {@code directory}: two independent {@code WHERE directory = ? ORDER BY
* time_updated DESC LIMIT 1} queries (one for the id, one for the model) can each pick a
* DIFFERENT row once more than one session shares a directory (fleetd #234 — the default
* no-worktree spawn shares the lead's cwd with every other worker and every past session ever
* run there), silently comparing a profile's requested model against a session that is not even
* the one whose id was returned. Keying on {@code id} instead makes that impossible: the model
* read back is always the SAME session {@link #sessionIdForDirectory} (or a cached copy of its
* answer) already resolved.
*
* <p>Kept as its own query and its own connection, independent from {@link
* #sessionIdForDirectory}: a database whose schema predates the {@code model} column (or any
* other read failure on this column alone) can never take id resolution down with it. That
* would be a regression of the id-resolution feature #209 shipped; this method degrades on its
* own.
*
* <p>Never throws, and every failure mode — no matching row, a missing/unreadable database, a
* missing {@code model} column, a null/blank {@code model} value, or JSON that does not parse
* into {@code {"id": "...", "providerID": "..."}} with a non-blank {@code id} — resolves to
* {@code null}. That is UNKNOWN evidence, not a mismatch signal: the caller must never
* quarantine a profile on the strength of a {@code null} here.
* quarantine a profile on the strength of a {@code null} here. A blank/null {@code sessionId}
* (the id is not resolved yet) is UNKNOWN too, for the same reason — never call this with one.
*
* @param directory the worker's cwd, as resolved for this spawn
* @param sessionId the session id already resolved by {@link #sessionIdForDirectory} for this
* spawn — never re-derived from {@code directory} here
* @return the actual model, or {@code null} when unknown
*/
ActualModel actualModelForDirectory(String directory) {
if (directory == null || directory.isBlank()) {
ActualModel actualModelForSessionId(String sessionId) {
if (sessionId == null || sessionId.isBlank()) {
return null;
}
if (!Files.isRegularFile(databasePath)) {
@@ -159,10 +187,10 @@ final class OpenCodeSessionDiscovery {
// exact condition — do not double-log it here.
return null;
}
String sql = "SELECT model FROM session WHERE directory = ? ORDER BY time_updated DESC LIMIT 1";
String sql = "SELECT model FROM session WHERE id = ?";
try (Connection connection = openReadOnly();
PreparedStatement statement = connection.prepareStatement(sql)) {
statement.setString(1, directory);
statement.setString(1, sessionId);
try (ResultSet rows = statement.executeQuery()) {
if (rows.next()) {
return parseModel(rows.getString("model"));
@@ -746,7 +746,7 @@ public final class MessageService {
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
TurnToken token = new TurnToken(target, reply);
TurnToken token = new TurnToken(target, reply, content);
// The send has won the lock; the accepted-delivery hook records delegator ownership
// here (CB-548). It runs BEFORE enqueue so a throwing hook — onAccepted is now a
// public callback — fails the send without queuing a message that would orphan.
@@ -9,8 +9,10 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ScheduledExecutorService;
@@ -68,6 +70,12 @@ public final class ReplyPushLoop {
static final String QUESTIONS_NUDGE_FORMAT =
"%d workers are paused on a question — run fleet_poll(ticket=...) for each, then answer "
+ "with fleet_send(turnId=..., content=...): %s";
static final String BACKEND_INCIDENT_NUDGE_FORMAT =
"Backend credential %s is cooling for %d remaining seconds; affected profiles: %s; "
+ "affected workers: %s. Run fleet_list to see more.";
static final String BACKEND_TARGET_UNMAPPED_NUDGE_FORMAT =
"A backend error on worker %s could not be mapped to a credential (%s) — no cool-off "
+ "was applied. Run fleet_list to check that worker.";
private final PrimaryRegistry primaryRegistry;
private final AgentControl agents;
@@ -93,6 +101,13 @@ public final class ReplyPushLoop {
* how depleted an older, still-open question's count is.
*/
private final ConcurrentHashMap<String, PendingQuestion> pendingQuestions = new ConcurrentHashMap<>();
/** Backend incidents awaiting one successful delivery, keyed by incident id and owning lead. */
private final ConcurrentHashMap<IncidentLead, PendingIncident> pendingIncidents = new ConcurrentHashMap<>();
/** Incident/lead pairs already delivered. They make repeated incident reports one-shot. */
private final Set<IncidentLead> deliveredIncidents = ConcurrentHashMap.newKeySet();
private final ConcurrentHashMap<UnmappedTargetLead, PendingUnmappedTarget> pendingUnmappedTargets =
new ConcurrentHashMap<>();
private final Set<UnmappedTargetLead> deliveredUnmappedTargets = ConcurrentHashMap.newKeySet();
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets and/or questions). */
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
@@ -178,7 +193,20 @@ public final class ReplyPushLoop {
* tracked per question, not per lead per source).
*/
private record PendingQuestion(String turnId, String ticket, String target, String lead,
String question, int nudgeCount) {
String question, int nudgeCount) {
}
private record IncidentLead(String incidentId, String lead) {
}
private record PendingIncident(IncidentLead key, String credential, List<String> profiles,
List<String> targets, int remainingCoolOffSeconds, int nudgeCount) {
}
private record UnmappedTargetLead(String target, String reason, String lead) {
}
private record PendingUnmappedTarget(UnmappedTargetLead key, int nudgeCount) {
}
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
@@ -192,6 +220,24 @@ public final class ReplyPushLoop {
.collect(Collectors.toUnmodifiableSet());
}
private List<PendingIncident> pendingIncidentsFor(String lead) {
return pendingIncidents.values().stream().filter(i -> lead.equals(i.key().lead())).toList();
}
private Set<IncidentLead> pendingIncidentKeysFor(String lead) {
return pendingIncidentsFor(lead).stream().map(PendingIncident::key)
.collect(Collectors.toUnmodifiableSet());
}
private List<PendingUnmappedTarget> pendingUnmappedTargetsFor(String lead) {
return pendingUnmappedTargets.values().stream().filter(i -> lead.equals(i.key().lead())).toList();
}
private Set<UnmappedTargetLead> pendingUnmappedTargetKeysFor(String lead) {
return pendingUnmappedTargetsFor(lead).stream().map(PendingUnmappedTarget::key)
.collect(Collectors.toUnmodifiableSet());
}
/**
* The reply-source reminder count {@link #decide} should see for {@code lead} on this tick:
* the <em>minimum</em> nudge count among the reply targets currently pending for it (CB-598).
@@ -236,6 +282,15 @@ public final class ReplyPushLoop {
return min == Integer.MAX_VALUE ? 0 : min;
}
private int minIncidentNudgeCountFor(String lead) {
return pendingIncidentsFor(lead).stream().mapToInt(PendingIncident::nudgeCount).min().orElse(0);
}
private int minUnmappedTargetNudgeCountFor(String lead) {
return pendingUnmappedTargetsFor(lead).stream().mapToInt(PendingUnmappedTarget::nudgeCount)
.min().orElse(0);
}
/**
* Pure decision function: examine everything pending for {@code lead} — reply targets and
* tickets alike — and return what the loop should do.
@@ -272,17 +327,35 @@ public final class ReplyPushLoop {
* @param questionReminderCount the lowest nudge count among questions open for this lead
*/
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount) {
return decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
minIncidentNudgeCountFor(lead));
}
/** As above, with backend incidents as a fourth, independently bounded source. */
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount,
int incidentReminderCount) {
return decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
incidentReminderCount, minUnmappedTargetNudgeCountFor(lead));
}
/** As above, with unmapped backend targets as a fifth, independently bounded source. */
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount,
int incidentReminderCount, int unmappedTargetReminderCount) {
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
boolean hasQuestionWork = !pendingQuestionTurnIdsFor(lead).isEmpty();
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork) {
boolean hasIncidentWork = !pendingIncidentKeysFor(lead).isEmpty();
boolean hasUnmappedTargetWork = !pendingUnmappedTargetKeysFor(lead).isEmpty();
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork && !hasUnmappedTargetWork) {
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
return Action.STOP;
}
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
boolean questionEligible = hasQuestionWork && questionReminderCount < maxReminders;
if (!replyEligible && !ticketEligible && !questionEligible) {
boolean incidentEligible = hasIncidentWork && incidentReminderCount < maxReminders;
boolean unmappedTargetEligible = hasUnmappedTargetWork && unmappedTargetReminderCount < maxReminders;
if (!replyEligible && !ticketEligible && !questionEligible && !incidentEligible && !unmappedTargetEligible) {
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
maxReminders, lead);
countNudge("exhausted");
@@ -394,6 +467,48 @@ public final class ReplyPushLoop {
pendingQuestions.remove(turnId);
}
/**
* Queue a one-shot backend credential outage notice for every distinct lead that owns an
* affected worker. A successful injection records its {@code (incidentId, lead)} key, so a
* repeat report never reminds that lead again.
*/
public void onBackendIncident(String incidentId, Collection<String> targets, String credential,
Collection<String> profiles, int remainingCoolOffSeconds) {
Map<String, List<String>> targetsByLead = new ConcurrentHashMap<>();
for (String target : targets) {
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.warn("push: backend incident {} has no known lead for target {}", incidentId, target);
continue;
}
targetsByLead.computeIfAbsent(lead.get(), _ -> new ArrayList<>()).add(target);
}
List<String> profileNames = profiles.stream().sorted().toList();
for (var entry : targetsByLead.entrySet()) {
IncidentLead key = new IncidentLead(incidentId, entry.getKey());
if (deliveredIncidents.contains(key)) continue;
pendingIncidents.putIfAbsent(key, new PendingIncident(key, credential, profileNames,
entry.getValue().stream().sorted().toList(), remainingCoolOffSeconds, 0));
startOrCoalesce(entry.getKey());
}
}
/**
* Tell the owning lead that a classified backend error could not be tied to a credential.
* Without an owning lead, emit a warning because no control can act on the target.
*/
public void onBackendTargetUnmapped(String target, String reason) {
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.warn("push: backend target {} could not map to a credential: {}", target, reason);
return;
}
UnmappedTargetLead key = new UnmappedTargetLead(target, reason, lead.get());
if (deliveredUnmappedTargets.contains(key)) return;
pendingUnmappedTargets.putIfAbsent(key, new PendingUnmappedTarget(key, 0));
startOrCoalesce(lead.get());
}
// --- the schedule ----------------------------------------------------------------------------
/** Start a reminder schedule for {@code lead}, or join the one already running. */
@@ -425,18 +540,25 @@ public final class ReplyPushLoop {
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
Set<String> questionsBefore = pendingQuestionTurnIdsFor(lead);
Set<IncidentLead> incidentsBefore = pendingIncidentKeysFor(lead);
Set<UnmappedTargetLead> unmappedTargetsBefore = pendingUnmappedTargetKeysFor(lead);
int replyReminderCount = minReplyNudgeCountFor(lead);
int ticketReminderCount = minTicketNudgeCountFor(lead);
int questionReminderCount = minQuestionNudgeCountFor(lead);
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
int incidentReminderCount = minIncidentNudgeCountFor(lead);
int unmappedTargetReminderCount = minUnmappedTargetNudgeCountFor(lead);
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
incidentReminderCount, unmappedTargetReminderCount);
switch (action) {
case INJECT -> {
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
incidentReminderCount, unmappedTargetReminderCount);
scheduleNext(lead);
}
// Re-check after the configured backoff; the lead may become injectable soon.
case WAIT_BUSY -> scheduleNext(lead);
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore);
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, incidentsBefore,
unmappedTargetsBefore);
}
}
@@ -480,11 +602,25 @@ public final class ReplyPushLoop {
* decision-to-release window reclaims the schedule slot exactly like a raced-in reply or ticket.
*/
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
Set<String> questionsBefore) {
Set<String> questionsBefore) {
stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, pendingIncidentKeysFor(lead));
}
private void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
Set<String> questionsBefore, Set<IncidentLead> incidentsBefore) {
stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, incidentsBefore,
pendingUnmappedTargetKeysFor(lead));
}
private void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
Set<String> questionsBefore, Set<IncidentLead> incidentsBefore,
Set<UnmappedTargetLead> unmappedTargetsBefore) {
activeLeads.remove(lead);
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t))
|| pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t));
|| pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t))
|| pendingIncidentKeysFor(lead).stream().anyMatch(i -> !incidentsBefore.contains(i))
|| pendingUnmappedTargetKeysFor(lead).stream().anyMatch(i -> !unmappedTargetsBefore.contains(i));
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
scheduleNext(lead);
@@ -495,18 +631,22 @@ public final class ReplyPushLoop {
/** Send one combined nudge covering everything currently pending for {@code lead}. */
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount,
int questionReminderCount) {
int questionReminderCount, int incidentReminderCount,
int unmappedTargetReminderCount) {
// Re-read rather than threading it down from decide(): a reply can drain, a ticket be
// collected, or a question be answered (or another arrive), between the decision and the
// injection.
Set<String> replyTargets = pendingReplyTargetsFor(lead);
List<PendingTicket> tickets = pendingTicketsFor(lead);
List<PendingQuestion> questions = pendingQuestionsFor(lead);
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty()) {
List<PendingIncident> incidents = pendingIncidentsFor(lead);
List<PendingUnmappedTarget> unmappedTargets = pendingUnmappedTargetsFor(lead);
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty() && incidents.isEmpty()
&& unmappedTargets.isEmpty()) {
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
return;
}
String nudge = formatNudge(replyTargets, tickets, questions);
String nudge = formatNudge(replyTargets, tickets, questions, incidents, unmappedTargets);
try {
agents.send(lead, nudge);
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
@@ -515,6 +655,16 @@ public final class ReplyPushLoop {
questionReminderCount + 1, maxReminders,
replyTargets.size(), tickets.size(), questions.size());
countNudge("delivered");
for (PendingIncident incident : incidents) {
if (pendingIncidents.remove(incident.key(), incident)) {
deliveredIncidents.add(incident.key());
}
}
for (PendingUnmappedTarget unmappedTarget : unmappedTargets) {
if (pendingUnmappedTargets.remove(unmappedTarget.key(), unmappedTarget)) {
deliveredUnmappedTargets.add(unmappedTarget.key());
}
}
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
@@ -526,12 +676,13 @@ public final class ReplyPushLoop {
// item already at or over the cap keeps riding along in the text (still pending, still
// named) but its extra bumps here are inert: decide() already treats it as ineligible once
// its count reaches maxReminders.
bumpNudgeCounts(replyTargets, tickets, questions);
bumpNudgeCounts(replyTargets, tickets, questions, incidents, unmappedTargets);
}
/** Record that every one of these items was just named in a sent (or attempted) nudge. */
private void bumpNudgeCounts(Set<String> replyTargets, List<PendingTicket> tickets,
List<PendingQuestion> questions) {
List<PendingQuestion> questions, List<PendingIncident> incidents,
List<PendingUnmappedTarget> unmappedTargets) {
for (String target : replyTargets) {
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
}
@@ -544,6 +695,14 @@ public final class ReplyPushLoop {
new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(),
e.nudgeCount() + 1));
}
for (PendingIncident incident : incidents) {
pendingIncidents.computeIfPresent(incident.key(), (id, e) -> new PendingIncident(e.key(),
e.credential(), e.profiles(), e.targets(), e.remainingCoolOffSeconds(), e.nudgeCount() + 1));
}
for (PendingUnmappedTarget unmappedTarget : unmappedTargets) {
pendingUnmappedTargets.computeIfPresent(unmappedTarget.key(), (id, e) ->
new PendingUnmappedTarget(e.key(), e.nudgeCount() + 1));
}
}
/** Schedule the next tick on the scheduler thread pool. */
@@ -556,7 +715,8 @@ public final class ReplyPushLoop {
/** Render everything pending for one lead as a single nudge line. */
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets,
List<PendingQuestion> questions) {
List<PendingQuestion> questions, List<PendingIncident> incidents,
List<PendingUnmappedTarget> unmappedTargets) {
List<String> parts = new ArrayList<>();
if (!replyTargets.isEmpty()) {
parts.add(formatRepliesNudge(replyTargets));
@@ -567,6 +727,12 @@ public final class ReplyPushLoop {
if (!questions.isEmpty()) {
parts.add(formatQuestionsNudge(questions));
}
if (!incidents.isEmpty()) {
parts.addAll(incidents.stream().map(ReplyPushLoop::formatBackendIncidentNudge).toList());
}
if (!unmappedTargets.isEmpty()) {
parts.addAll(unmappedTargets.stream().map(ReplyPushLoop::formatUnmappedTargetNudge).toList());
}
return String.join(" | ", parts);
}
@@ -606,6 +772,16 @@ public final class ReplyPushLoop {
return QUESTIONS_NUDGE_FORMAT.formatted(pending.size(), ids);
}
private static String formatBackendIncidentNudge(PendingIncident incident) {
return BACKEND_INCIDENT_NUDGE_FORMAT.formatted(incident.credential(), incident.remainingCoolOffSeconds(),
String.join(", ", incident.profiles()), String.join(", ", incident.targets()));
}
private static String formatUnmappedTargetNudge(PendingUnmappedTarget unmappedTarget) {
return BACKEND_TARGET_UNMAPPED_NUDGE_FORMAT.formatted(unmappedTarget.key().target(),
unmappedTarget.key().reason());
}
// --- lifecycle -----------------------------------------------------------------------------
/**
@@ -627,6 +803,10 @@ public final class ReplyPushLoop {
pendingReplies.clear();
pendingTickets.clear();
pendingQuestions.clear();
pendingIncidents.clear();
deliveredIncidents.clear();
pendingUnmappedTargets.clear();
deliveredUnmappedTargets.clear();
}
/** @see #stop() */
@@ -9,12 +9,19 @@ import java.util.concurrent.CompletableFuture;
public final class TurnToken {
private final String target;
private final CompletableFuture<Rendezvous.Resolution> waiter;
private final String injectedText;
public TurnToken(String target, CompletableFuture<Rendezvous.Resolution> waiter) {
this(target, waiter, null);
}
public TurnToken(String target, CompletableFuture<Rendezvous.Resolution> waiter, String injectedText) {
this.target = target;
this.waiter = waiter;
this.injectedText = injectedText;
}
public String target() { return target; }
public CompletableFuture<Rendezvous.Resolution> waiter() { return waiter; }
public String injectedText() { return injectedText; }
}
@@ -0,0 +1,200 @@
package dev.ltms.fleet.placement;
import java.util.ArrayList;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.OptionalLong;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.LongSupplier;
/**
* Fleetd #201 / #227: correlates classified backend errors (credential outage or provider 5xx,
* produced elsewhere by the classifier that turns raw pane text into a typed event — this class
* knows nothing about pane text, profiles, sessions, launchers, or leads) and decides, purely from
* counts and timing, when a credential's backend is out.
*
* <p>Two classified errors on the <b>same credential</b> — never the profile name, never the error
* text, see {@code FleetConfig.Profile#effectiveCredentialId()} like {@link BackendQuarantine} —
* <b>from two distinct targets</b> inside a 60-second window is treated as an outage: {@link
* #record} then returns an {@link Incident} and starts a 60-second cool-off for that credential. The
* threshold counts distinct targets, not raw events, on purpose: the classifier is a heuristic and a
* valid member report can quote an {@code API Error:} line, so one target repeating that line twice
* must never remove fleet capacity by itself. Two different targets independently producing a
* classified error is far less likely to be a coincidence, and a real backend outage hits every
* target on that credential anyway, so this loses nothing against the case being protected against.
* This is deliberately a different store from {@link BackendQuarantine}: that one holds a
* 1800-second exhaustion cooldown for a spent credential, and reusing it here would both use the
* wrong duration and report "backend exhausted" for what is a short transient fault.
*
* <p>Correlation and cool-off live in <b>one class</b> so that crossing the threshold and setting
* the cool-off deadline happen as a single atomic update: every {@link #record} call goes through
* {@link ConcurrentHashMap#compute}, which serializes on the credential's map bucket, so a
* concurrent second and third event on the same credential can never both observe "about to cross"
* and both mint an incident.
*
* <p>The clock is injected ({@link LongSupplier}, conventionally {@code System::nanoTime} like
* {@link BackendQuarantine}), never read inline, so the window and cool-off are testable without a
* real sleep.
*/
public final class BackendOutagePolicy {
/** Distinct targets a credential needs a classified error from to declare an outage. */
public static final int THRESHOLD = 2;
/** How long a credential's evidence stays fresh, inclusive of both endpoints. */
public static final long WINDOW_NANOS = 60_000_000_000L;
/** How long a credential sits out once the threshold is crossed. */
public static final long COOLOFF_NANOS = 60_000_000_000L;
private final ConcurrentHashMap<String, CredentialState> states = new ConcurrentHashMap<>();
private final LongSupplier nowNanos;
private final AtomicLong incidentSequence = new AtomicLong();
public BackendOutagePolicy(LongSupplier nowNanos) {
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
}
/**
* Records one classified backend error for {@code credentialId} against {@code target} (e.g. a
* session or worker id — this class never interprets it, only collects it for the incident's
* affected-targets list, and counts distinct targets toward the threshold) with {@code reason}
* (the classifier's free-text reason, kept per event for whoever renders the eventual notice —
* every event's reason is kept even when the same target repeats, so {@link Incident#reasons()}
* can be longer than {@link Incident#evidenceCount()}).
*
* <p>Returns a populated {@link Incident} only at the exact moment a <b>second distinct target</b>
* is seen for this credential inside the window — never before, and never again while the
* resulting cool-off is active. A repeat error from a target already counted does not advance the
* threshold. While a credential is cooling off, a fresh error is ignored outright: it neither
* extends the deadline nor produces another incident. Once the cool-off has elapsed, the next
* error clears the old evidence and starts a brand-new window — two fresh, distinct targets are
* required to rearm.
*/
public Optional<Incident> record(String credentialId, String target, String reason) {
Objects.requireNonNull(credentialId, "credentialId");
Objects.requireNonNull(target, "target");
Objects.requireNonNull(reason, "reason");
long now = nowNanos.getAsLong();
AtomicReference<Incident> minted = new AtomicReference<>();
states.compute(credentialId, (_, existing) -> {
CredentialState state = existing;
if (state != null && state.inCoolOff()) {
if (now < state.coolOffUntilNanos) {
return state; // still cooling off: no extension, no incident
}
state = null; // cool-off elapsed: evidence is cleared, rearm from scratch
}
if (state == null || now - state.firstEventNanos > WINDOW_NANOS) {
return CredentialState.first(now, target, reason);
}
CredentialState advanced = state.withAdditionalEvidence(target, reason);
if (advanced.evidenceCount() < THRESHOLD) {
return advanced;
}
long coolOffUntilNanos = now + COOLOFF_NANOS;
minted.set(new Incident(
"outage-" + credentialId + "-" + incidentSequence.incrementAndGet(),
credentialId,
Set.copyOf(advanced.targets),
advanced.evidenceCount(),
toSecondsRoundedUp(WINDOW_NANOS),
toSecondsRoundedUp(COOLOFF_NANOS),
List.copyOf(advanced.reasons)));
return advanced.enteringCoolOff(coolOffUntilNanos);
});
return Optional.ofNullable(minted.get());
}
/** Seconds left on {@code credentialId}'s cool-off, or empty when it is not cooling off. */
public OptionalLong remainingCoolOffSeconds(String credentialId) {
Objects.requireNonNull(credentialId, "credentialId");
CredentialState state = states.get(credentialId);
if (state == null || !state.inCoolOff()) {
return OptionalLong.empty();
}
long remaining = state.coolOffUntilNanos - nowNanos.getAsLong();
return remaining > 0 ? OptionalLong.of(toSecondsRoundedUp(remaining)) : OptionalLong.empty();
}
private static long toSecondsRoundedUp(long nanos) {
return (nanos + 999_999_999L) / 1_000_000_000L;
}
/**
* One credential crossing the outage threshold. {@code remainingCoolOffSeconds} is the cool-off
* length as observed at the moment of minting — this incident is only ever produced right as the
* cool-off starts, so it is always the full {@link #COOLOFF_NANOS} rounded up.
*
* <p>{@code evidenceCount} is {@code targets.size()} — the threshold is on distinct targets, not
* raw events — while {@code reasons} keeps every event's reason, including repeats from a target
* already counted. The two are deliberately different lengths: a single target hammering the same
* classified error never grows {@code evidenceCount} past 1, but each occurrence still lands in
* {@code reasons} for whoever renders the notice.
*/
public record Incident(
String id,
String credentialId,
Set<String> targets,
int evidenceCount,
long windowSeconds,
long remainingCoolOffSeconds,
List<String> reasons) {
}
/** Evidence accumulated for one credential since its window opened, or its active cool-off. */
private static final class CredentialState {
final long firstEventNanos;
final Set<String> targets;
final List<String> reasons;
final long coolOffUntilNanos; // 0 means "not cooling off"
private CredentialState(long firstEventNanos, Set<String> targets, List<String> reasons,
long coolOffUntilNanos) {
this.firstEventNanos = firstEventNanos;
this.targets = targets;
this.reasons = reasons;
this.coolOffUntilNanos = coolOffUntilNanos;
}
static CredentialState first(long nowNanos, String target, String reason) {
Set<String> targets = new LinkedHashSet<>();
targets.add(target);
List<String> reasons = new ArrayList<>();
reasons.add(reason);
return new CredentialState(nowNanos, targets, reasons, 0L);
}
CredentialState withAdditionalEvidence(String target, String reason) {
Set<String> newTargets = new LinkedHashSet<>(targets);
newTargets.add(target);
List<String> newReasons = new ArrayList<>(reasons);
newReasons.add(reason);
return new CredentialState(firstEventNanos, newTargets, newReasons, 0L);
}
CredentialState enteringCoolOff(long coolOffUntilNanos) {
return new CredentialState(firstEventNanos, targets, reasons, coolOffUntilNanos);
}
/** Distinct targets seen so far — the threshold counts this, never {@code reasons.size()}. */
int evidenceCount() {
return targets.size();
}
boolean inCoolOff() {
return coolOffUntilNanos > 0;
}
}
}
@@ -46,7 +46,8 @@ public record MemberSession(
String worktree,
String branch,
CharterReceipt charterReceipt,
String agentSessionId) {
String agentSessionId,
String failureReason) {
/** One-shot worker lifecycle states. */
public enum State {
@@ -54,6 +55,7 @@ public record MemberSession(
READY,
BUSY,
DONE,
BACKEND_ERROR,
FAILED,
RELEASED
}
@@ -69,25 +71,35 @@ public record MemberSession(
long lastActivityAtNanos, int turnCount, State state,
String worktree, String branch) {
this(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, null, null);
lastActivityAtNanos, turnCount, state, worktree, branch, null, null, null);
}
/** Backward-compatible shape without a backend failure reason. */
public MemberSession(String paneId, String terminalId, String profile, MemberRole role,
String cwd, String ownerTerminal, long spawnedAtNanos,
long lastActivityAtNanos, int turnCount, State state,
String worktree, String branch, CharterReceipt charterReceipt,
String agentSessionId) {
this(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, null);
}
/** Return a copy of this session in {@code state}. */
public MemberSession withState(State state) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
/** Return a copy with {@code lastActivityAtNanos} updated to {@code nowNanos}. */
public MemberSession withActivity(long nowNanos) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
nowNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
nowNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
/** Return a copy with the turn count incremented and activity timestamped at {@code nowNanos}. */
public MemberSession bumpTurn(long nowNanos) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId);
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
/**
@@ -98,6 +110,12 @@ public record MemberSession(
*/
public MemberSession withAgentSessionId(String agentSessionId) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
/** Return a copy with the durable backend failure detail. */
public MemberSession withFailureReason(String failureReason) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
}
@@ -658,6 +658,9 @@ public final class SessionManager implements TurnListener {
m.put("profile", session.profile());
m.put("role", session.role() == null ? "dev" : session.role().wireName());
m.put("state", session.state().name().toLowerCase());
if (session.failureReason() != null) {
m.put("failureReason", session.failureReason());
}
if (session.worktree() != null) {
m.put("worktree", session.worktree());
}
@@ -718,6 +721,36 @@ public final class SessionManager implements TurnListener {
completeTurn(target, false);
}
/**
* Record that a member backend failed after a delegated ticket expired. This CAS loop accepts
* either side of the completion race: {@code BUSY -> BACKEND_ERROR} or
* {@code DONE -> BACKEND_ERROR}. A released or unknown session is never recreated.
*
* @return {@code true} when the error is recorded on a live session
*/
public boolean onBackendError(String target, String reason) {
while (true) {
MemberSession current = findByTerminal(target);
if (current == null) {
log.warn("backend error for unknown member terminal={}", target);
return false;
}
if (current.state() == MemberSession.State.RELEASED) {
return false;
}
if (current.state() == MemberSession.State.BACKEND_ERROR) {
return true;
}
MemberSession updated = current.withState(MemberSession.State.BACKEND_ERROR)
.withFailureReason(reason).withActivity(nowNanos.getAsLong());
if (replace(current, updated)) {
log.warn("member terminal={} pane={} transitioned {} -> BACKEND_ERROR: {}",
target, current.paneId(), current.state(), reason);
return true;
}
}
}
@Override
public boolean hasPostTurnAction(String target) {
if (!clearAfterTurn) return false;
@@ -736,10 +769,9 @@ public final class SessionManager implements TurnListener {
if (current == null || current.state() != MemberSession.State.BUSY) return false;
long now = nowNanos.getAsLong();
MemberSession updated = current.withState(MemberSession.State.DONE).withActivity(now);
if (replace(current, updated)) {
log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}",
target, current.paneId(), updated.turnCount());
}
if (!replace(current, updated)) return false;
log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}",
target, current.paneId(), updated.turnCount());
if (contextCap > 0 && updated.turnCount() >= contextCap) {
release(current.paneId());
return false;
@@ -1,8 +1,12 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
@@ -26,19 +30,66 @@ class RequiredSecretEnvVarsTest {
return FleetConfig.load(f);
}
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
}
@Test
void collectsATokenEnvPerNonSubscriptionProfile(@TempDir Path dir) throws Exception {
void explicitlyConfiguredUnsetTokenEnvWarnsWithCurrentWording(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
tokenEnv: AI_GATEWAY_TOKEN
tokenEnv: CB115_MISSING_TOKEN_ENV_819367
""");
Map<String, List<String>> required = Fleetd.requiredSecretEnvVars(cfg);
assertTrue(required.containsKey("AI_GATEWAY_TOKEN"));
assertEquals(List.of("profile 'local' tokenEnv"), required.get("AI_GATEWAY_TOKEN"));
assertTrue(required.containsKey("CB115_MISSING_TOKEN_ENV_819367"));
assertEquals(List.of("profile 'local' tokenEnv"), required.get("CB115_MISSING_TOKEN_ENV_819367"));
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRequiredSecrets(cfg);
} finally {
detach(appender);
}
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == ch.qos.logback.classic.Level.WARN
&& e.getFormattedMessage().contains("startup secret CB115_MISSING_TOKEN_ENV_819367: MISSING")
&& e.getFormattedMessage().contains("profile 'local' tokenEnv")),
"an explicitly configured but unset tokenEnv must retain the startup warning");
}
@Test
void profileWithoutTokenEnvDoesNotRequireTheOldDefault(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
""");
assertTrue(Fleetd.requiredSecretEnvVars(cfg).isEmpty(),
"an absent tokenEnv means this profile needs no token, not FLEETD_WORKER_TOKEN");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRequiredSecrets(cfg);
} finally {
detach(appender);
}
assertFalse(appender.list.stream().anyMatch(e -> e.getLevel() == ch.qos.logback.classic.Level.WARN),
"a profile without tokenEnv must not produce a startup secret warning");
}
@Test
@@ -196,6 +196,26 @@ class CompletionResolverTest {
assertEquals("complete report", waiter.getNow(null).text());
}
@Test
void suppressesABareEchoWithOnlyTuiChrome() {
String injected = "Load the implementer skill. You own fleetd #999. ".repeat(12);
String scrape = injected + "\nDev auto - GPT-5.6 Terra OpenAI";
assertTrue(CompletionResolver.echoesInjectedBrief(scrape, injected));
}
@Test
void pinsTheMaximumTuiChromeExcess() {
String injected = "a".repeat(CompletionResolver.ECHO_MIN_CHARS);
String underMargin = injected + "b".repeat(CompletionResolver.MAX_ECHO_EXCESS_CHARS);
String overMargin = injected + "b".repeat(CompletionResolver.MAX_ECHO_EXCESS_CHARS + 1);
assertTrue(CompletionResolver.echoesInjectedBrief(underMargin, injected),
"the configured excess itself remains an echoed brief");
assertFalse(CompletionResolver.echoesInjectedBrief(overMargin, injected),
"one character beyond the excess must preserve the scrape as a real report");
}
@Test
void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
@@ -502,7 +522,7 @@ class CompletionResolverTest {
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
@@ -523,7 +543,7 @@ class CompletionResolverTest {
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
@@ -651,7 +671,7 @@ class CompletionResolverTest {
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
@@ -704,7 +724,7 @@ class CompletionResolverTest {
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
@@ -731,7 +751,7 @@ class CompletionResolverTest {
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
@@ -871,7 +891,7 @@ class CompletionResolverTest {
BackendErrorPatternLookup backendErrors = target -> Pattern.compile("(?i)usage limit");
java.util.List<String> exhaustedNotified = new java.util.ArrayList<>();
java.util.List<String> backendErrorNotified = new java.util.ArrayList<>();
ExhaustionSink exhaustionSink = (target, reason) -> exhaustedNotified.add(target);
ExhaustionSink exhaustionSink = (target, reason, profile) -> exhaustedNotified.add(target);
BackendErrorSink backendErrorSink = (target, matchedLine, reason) -> backendErrorNotified.add(target);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
exhausted, exhaustionSink, backendErrors, backendErrorSink);
@@ -0,0 +1,41 @@
package dev.ltms.fleet.inject;
import org.junit.jupiter.api.Test;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #234. Round 3: calls the REAL production factory, {@link ExhaustionSink#forwardingTo},
* rather than rebuilding its shape locally — a round-2 version of this test built its own copy of
* the forwarder inline, so mutating {@code Fleetd.java}'s actual forwarder back into a lambda left
* that copy untouched and the test kept passing while production had regressed. Calling the shared
* factory here means the same object under test IS the object {@code Fleetd.java} builds via
* {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)}.
*
* <p>Round 4: the 3-arg overload is now {@link ExhaustionSink}'s single abstract method, so {@code
* forwardingTo} itself is a one-line lambda and every sink below can safely be one too — there is
* no 2-arg overload left for any of them to silently bind to instead. This test still catches a
* regression inside {@code forwardingTo}'s body (e.g. one that calls the 2-arg default and drops
* the hint that way) because it still goes through the shared factory rather than a rebuilt copy.
*/
class ExhaustionSinkForwardingHazardTest {
@Test
void forwardingToDeliversTheProfileHintToWhateverSinkTheSupplierCurrentlyReturns() {
AtomicReference<String> hintSeenByRealSink = new AtomicReference<>("NEVER CALLED");
ExhaustionSink real = (target, reason, profileHint) -> hintSeenByRealSink.set(profileHint);
// Same shape Fleetd.java uses: a reference that starts at none() and is repointed later.
AtomicReference<ExhaustionSink> ref = new AtomicReference<>(ExhaustionSink.none());
ExhaustionSink forwarding = ExhaustionSink.forwardingTo(ref::get);
ref.set(real);
forwarding.onExhausted("term_x", "model mismatch", "gx");
assertEquals("gx", hintSeenByRealSink.get(),
"ExhaustionSink.forwardingTo must deliver the profile hint to the sink the supplier "
+ "currently returns — the exact object Fleetd.java's forwarder is built from");
}
}
@@ -583,6 +583,24 @@ class ClaudeCodeLauncherTest {
assertNull(env.get("GITEA_HOST"), "no forge host without a granted token");
}
@Test
void noTokenEnvSpawnsWithoutInjectingAnAuthToken() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"local-direct", "http://gx00.gw:8000", "coder", null, null,
List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), System::getenv);
assertDoesNotThrow(() -> {
launcher.spawn();
},
"a profile that needs no token must still spawn against its configured endpoint");
Map<String, String> env = startEnv(herdr);
assertFalse(env.containsKey("ANTHROPIC_AUTH_TOKEN"),
"an unconfigured tokenEnv must not inject an auth-token key");
}
// --- CB-117 orphan reap: the pure predicate --------------------------------
@Test
@@ -17,6 +17,7 @@ import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import org.junit.jupiter.api.Test;
@@ -34,6 +35,7 @@ import java.util.Optional;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
import static org.junit.jupiter.api.Assertions.*;
@@ -52,6 +54,19 @@ class OpenCodeLauncherTest {
null, null, gitTokenEnv, null, FleetConfig.Profile.KIND_OPENCODE);
}
/**
* A profile carrying an explicit {@code credentialId} (fleetd #234 defect 2) — distinct from
* the profile's own name, so a test can assert on the credential precisely rather than relying
* on {@code effectiveCredentialId()}'s profile-name fallback.
*/
private static FleetConfig.Profile opencodeCfgWithCredential(String profileName, String model,
String credentialId) {
return new FleetConfig.Profile(profileName, null, model, null, "FLEETD_WORKER_TOKEN",
List.of("opencode"), "tab", "fleetd-workers", "opencode: {model} #{n}", null,
null, List.of(), null, null, FleetConfig.Profile.KIND_OPENCODE, Map.of(), 1.0f,
null, false, null, credentialId, null);
}
/** Gate-disabled launcher whose per-spawn config dirs land under an inspectable temp root. */
private static OpenCodeLauncher service(FakeHerdr herdr, Path configRoot, FleetConfig.Profile cfg) {
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
@@ -910,7 +925,7 @@ class OpenCodeLauncherTest {
// credentialId — the profile that actually escaped the fleet's accounting.
FleetConfig.Profile cfg = opencodeCfg("opencode/nemotron-3-ultra-free", null, null);
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(target + "|" + reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(target + "|" + reason);
OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, sink);
SessionManager sessions = new SessionManager(launcher);
@@ -945,7 +960,7 @@ class OpenCodeLauncherTest {
void aProviderPrefixedModelMatchingBothIdAndProviderIsNotAMismatch(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null);
PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
.spawn(new SpawnRequest(null, "/work/dir", null));
@@ -961,7 +976,7 @@ class OpenCodeLauncherTest {
void aGxProviderPrefixedModelMatchingBothIdAndProviderIsNotAMismatch(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
FleetConfig.Profile cfg = opencodeCfg("gx/deepseek-v4-flash", null, null);
PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
.spawn(new SpawnRequest(null, "/work/dir", null));
@@ -983,7 +998,7 @@ class OpenCodeLauncherTest {
void aMissingProviderIdInTheEvidenceIsUnknownNotAMismatchWhenTheIdMatches(
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null);
PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
.spawn(new SpawnRequest(null, "/work/dir", null));
@@ -1004,7 +1019,7 @@ class OpenCodeLauncherTest {
void aMissingProviderIdInTheEvidenceStillCatchesARealIdMismatch(
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null);
PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
.spawn(new SpawnRequest(null, "/work/dir", null));
@@ -1028,7 +1043,7 @@ class OpenCodeLauncherTest {
void aBareModelWithNoProviderPrefixMatchesOnIdAloneAndIsNotAMismatch(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
FleetConfig.Profile cfg = opencodeCfg("deepseek-v4-flash", null, null);
PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
.spawn(new SpawnRequest(null, "/work/dir", null));
@@ -1050,7 +1065,7 @@ class OpenCodeLauncherTest {
void aRealIdMismatchLogsAnErrorNamingBothModelsAndQuarantinesThroughTheSink(
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(target + "|" + reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(target + "|" + reason);
FleetConfig.Profile cfg = opencodeCfg("opencode/nemotron-3-ultra-free", null, null);
Logger logger = (Logger) LoggerFactory.getLogger(OpenCodeLauncher.class);
@@ -1097,7 +1112,7 @@ class OpenCodeLauncherTest {
void unknownOrUnparseableModelEvidenceNeverQuarantines(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null);
PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
.spawn(new SpawnRequest(null, "/work/dir", null));
@@ -1123,7 +1138,7 @@ class OpenCodeLauncherTest {
void aProfileWithNoConfiguredModelIsNeverCheckedForAMismatch(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason) -> exhausted.add(reason);
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
FleetConfig.Profile cfg = opencodeCfg(null, null, null);
PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
.spawn(new SpawnRequest(null, "/work/dir", null));
@@ -1133,4 +1148,200 @@ class OpenCodeLauncherTest {
assertEquals("ses_x", handle.agentSessionId());
assertTrue(exhausted.isEmpty(), "no model configured → nothing to compare: " + exhausted);
}
// --- fleetd #234 defect 1: key the model check on the RESOLVED id, not the shared directory --
/**
* The exact shape fleetd #234 reported: {@code fleet_spawn} with no {@code worktree:} shares
* the lead's cwd across every worker, so more than one session row can exist for the SAME
* {@code directory}. Once THIS handle's own session id is resolved, a sibling member spawned
* later into the same shared directory — writing a NEWER, unrelated row — must never make the
* already-resolved session look mismatched. Today's code re-derives "the newest row in this
* directory" on every call (both for the id AND, independently, for the model), so it would
* pick up the sibling's row on the second call and flag a false mismatch AND flip the returned
* id. The fix (fleetd #234) makes the id sticky once resolved and reads the model back for
* exactly that id (see {@link OpenCodeSessionDiscovery#actualModelForSessionId}) — never
* "whatever is newest in the directory right now."
*/
@Test
void modelCheckReadsTheResolvedSessionsOwnRowNotWhateverIsNewestInTheSharedDirectory(
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
List<String> exhausted = new ArrayList<>();
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null);
PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
.spawn(new SpawnRequest(null, "/work/dir", null));
// Our own session's row, correctly matching the profile's requested model.
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_ours", "/work/dir", 1000L,
"{\"id\":\"gpt-5.6-terra\",\"providerID\":\"openai\"}");
assertEquals("ses_ours", handle.agentSessionId(), "resolves to our own session");
assertTrue(exhausted.isEmpty(), "matching model → no mismatch on first resolve: " + exhausted);
// A sibling member, spawned later into the SAME shared directory (no worktree, fleetd
// #234's default), writes a newer row running a totally different model.
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_sibling", "/work/dir", 9000L,
"{\"id\":\"deepseek-v4-flash\",\"providerID\":\"gx\"}");
assertEquals("ses_ours", handle.agentSessionId(),
"the session id, once resolved, must not flip to a sibling sharing the directory");
assertTrue(exhausted.isEmpty(),
"a sibling's later, unrelated row in the same shared directory must never be read "
+ "as OUR session's model: " + exhausted);
}
// --- fleetd #234 defect 2: the quarantine the ERROR announces must actually happen -----------
/**
* fleetd #234: {@link OpenCodeLauncher.SessionAwareHandle#checkModelMatch} fires from {@code
* agentSessionId()}, which {@code SessionManager.acquire()} calls to build the very first
* {@code MemberSession} record — BEFORE that session is put into the registry {@code
* sessions.roster()} reads. A sink that resolves {@code target -> profile} ONLY through the
* roster (today's {@code Fleetd.java} code, before this fix) therefore finds nothing at this
* exact moment and silently does not quarantine, even though it just logged an ERROR saying it
* would. This test drives the REAL path — {@code SessionManager.acquire()} — not the sink
* directly, because the bug is entirely about this ordering; a direct-sink test cannot see it
* (and is exactly why #175's own test suite, which only ever called the sink directly or after
* registration, never caught this).
*
* <p>The sink under test mirrors {@code Fleetd.main()}'s real wiring after the fix: resolve
* via the roster first (unchanged for {@code CompletionResolver}'s two call sites), falling
* back to the profile hint {@link OpenCodeLauncher} now supplies via {@link
* ExhaustionSink#onExhausted(String, String, String)} when the roster lookup misses.
*/
@Test
void aSpawnTimeModelMismatchActuallyQuarantinesTheCredentialThroughTheRealAcquirePath(
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
FleetConfig.Profile cfg = opencodeCfgWithCredential(
"terra", "opencode/nemotron-3-ultra-free", "openai-shared");
Map<String, FleetConfig.Profile> profiles = Map.of(cfg.profile(), cfg);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800));
// fleetd #234, round 4: safely a lambda now — the 3-arg overload is the interface's single
// abstract method, so there is no 2-arg overload left to bind to instead.
ExhaustionSink sink = (target, reason, profileHint) -> {
FleetConfig.Profile profile = profileHint == null ? null : profiles.get(profileHint);
if (profile != null) {
quarantine.quarantine(profile.effectiveCredentialId());
}
};
FakeHerdr herdr = new FakeHerdr();
OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, sink);
SessionManager sessions = new SessionManager(launcher);
// The mismatching row exists BEFORE the spawn — reproducing fleetd #234's exact timing:
// opencode's session table already carries evidence by the moment acquire() first asks.
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
"{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}");
assertFalse(quarantine.isQuarantined("openai-shared"), "nothing quarantined before the spawn");
// The real production entrypoint: acquire() builds the MemberSession by calling
// handle.agentSessionId() BEFORE registry.put() runs.
MemberSession acquired = sessions.acquire(cfg.profile(), "/work/dir", null, null);
assertEquals("ses_x", acquired.agentSessionId(), "the id itself still resolves correctly");
assertTrue(quarantine.isQuarantined("openai-shared"),
"the mismatch fires DURING acquire(), before roster registration, and must still "
+ "reach the quarantine via the profile hint — not silently no-op");
}
/**
* The other half of the same proof: a sink that resolves {@code target -> profile} ONLY
* through the roster (i.e. ignores the profile hint entirely, ~today's pre-fix {@code
* Fleetd.java}) drops the SAME spawn-time mismatch silently — the credential is never
* quarantined even though {@link OpenCodeLauncher} logged the mismatch ERROR. This is the
* failure fleetd #234 reported, reproduced through the real {@code SessionManager.acquire()}
* path rather than asserted by inspecting the fix.
*/
@Test
void aRosterOnlySinkSilentlyDropsTheSpawnTimeQuarantine(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
FleetConfig.Profile cfg = opencodeCfgWithCredential(
"terra", "opencode/nemotron-3-ultra-free", "openai-shared");
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800));
// Deliberately ignores the profile hint — the pre-fix shape: only a roster lookup (modelled
// here as always empty, since acquire() has not registered the session yet either way).
ExhaustionSink rosterOnlySink = (target, reason, profile) -> { };
FakeHerdr herdr = new FakeHerdr();
OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, rosterOnlySink);
SessionManager sessions = new SessionManager(launcher);
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
"{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}");
MemberSession acquired = sessions.acquire(cfg.profile(), "/work/dir", null, null);
assertEquals("ses_x", acquired.agentSessionId(), "the id itself still resolves correctly");
assertFalse(quarantine.isQuarantined("openai-shared"),
"a roster-only sink cannot see this target yet — the quarantine silently never "
+ "happens, which is exactly fleetd #234 defect 2");
}
/**
* fleetd #234, round 2: the two tests above inject a sink DIRECTLY into {@link OpenCodeLauncher},
* which is not what {@code Fleetd.main} actually does. Production has one more hop: {@code
* Fleetd.java} builds the adapters (including {@code OpenCodeLauncher}) before {@code sessions}
* exists — a genuine construction-order cycle — so it hands the launcher a <em>forwarding</em>
* sink pointed at an {@code AtomicReference<ExhaustionSink>}, and only later {@code .set(...)}s
* that reference to the real sink once {@code sessions} is built.
*
* <p><strong>Round 3:</strong> a round-2 version of this test built its own copy of that
* forwarder's shape as an anonymous class. Mutating {@code Fleetd.java}'s ACTUAL forwarder back
* into a broken lambda left this test's own copy untouched, so it kept passing while production
* had regressed to the exact bug being fixed. This version instead calls {@link
* ExhaustionSink#forwardingTo}, the same factory {@code Fleetd.java} calls — the identical
* object, not a rebuilt copy of its shape — so a regression at either the {@code Fleetd.java}
* call site or inside {@code forwardingTo} itself has nowhere left to hide. See {@code
* ExhaustionSinkForwardingHazardTest} for the same factory exercised in isolation.
*/
@Test
void theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
FleetConfig.Profile cfg = opencodeCfgWithCredential(
"terra", "opencode/nemotron-3-ultra-free", "openai-shared");
Map<String, FleetConfig.Profile> profiles = Map.of(cfg.profile(), cfg);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800));
// Fleetd.java:176 — the forwarding sink is built BEFORE the real one can exist, and the
// launcher below is constructed against this forwarder, exactly like Fleetd.main. Calling
// the SAME factory Fleetd.java calls — ExhaustionSink.forwardingTo — rather than rebuilding
// the forwarder's shape here is the whole point (fleetd #234, round 3): a round-2 version of
// this test built its own copy, so mutating Fleetd.java's real forwarder back into a lambda
// left this test untouched. Routed through the shared factory, a regression at either the
// Fleetd.java call site (reverting to a hand-written lambda) or inside the factory body
// itself has nowhere left to hide from this test.
java.util.concurrent.atomic.AtomicReference<ExhaustionSink> exhaustionSinkRef =
new java.util.concurrent.atomic.AtomicReference<>(ExhaustionSink.none());
ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
FakeHerdr herdr = new FakeHerdr();
OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, forwardingExhaustionSink);
SessionManager sessions = new SessionManager(launcher);
// Fleetd.java:359-390 — the real sink is only built and pointed to AFTER `sessions` exists,
// same order as production. Safely a lambda (round 4): see the note on the forwarder above.
ExhaustionSink realSink = (target, reason, profileHint) -> {
FleetConfig.Profile profile = profileHint == null ? null : profiles.get(profileHint);
if (profile != null) {
quarantine.quarantine(profile.effectiveCredentialId());
}
};
exhaustionSinkRef.set(realSink);
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
"{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}");
assertFalse(quarantine.isQuarantined("openai-shared"), "nothing quarantined before the spawn");
MemberSession acquired = sessions.acquire(cfg.profile(), "/work/dir", null, null);
assertEquals("ses_x", acquired.agentSessionId(), "the id itself still resolves correctly");
assertTrue(quarantine.isQuarantined("openai-shared"),
"the profile hint must survive the Fleetd-style forwarding hop and reach the real "
+ "sink — a lambda forwarder drops it and this must go red");
}
}
@@ -146,14 +146,14 @@ class OpenCodeSessionDiscoveryTest {
"an unreadable database resolves to null, not an exception");
}
// --- fleetd #175: actualModelForDirectory / the model JSON column ---------------------------
// --- fleetd #175/#234: actualModelForSessionId / the model JSON column -----------------------
@Test
void parsesTheModelJsonIntoProviderAndId(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"id\":\"gpt-5.6-terra\",\"providerID\":\"openai\"}");
OpenCodeSessionDiscovery.ActualModel actual =
new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a");
new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa");
assertNotNull(actual, "a well-formed model JSON parses");
assertEquals("gpt-5.6-terra", actual.id());
@@ -161,28 +161,39 @@ class OpenCodeSessionDiscoveryTest {
}
@Test
void prefersTheModelOfTheMostRecentlyUpdatedRow(@TempDir Path root) throws Exception {
writeRecord(root, "ses_old", "/w/a", 1000L, "{\"id\":\"old-model\",\"providerID\":\"openai\"}");
writeRecord(root, "ses_new", "/w/a", 5000L, "{\"id\":\"new-model\",\"providerID\":\"openai\"}");
void looksUpByIdEvenWhenAnotherRowInTheSameDirectoryIsNewer(@TempDir Path root) throws Exception {
writeRecord(root, "ses_ours", "/work/dir", 1000L, "{\"id\":\"gpt-5.6-terra\",\"providerID\":\"openai\"}");
writeRecord(root, "ses_sibling", "/work/dir", 9000L, "{\"id\":\"deepseek-v4-flash\",\"providerID\":\"gx\"}");
assertEquals("new-model",
new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a").id(),
"the model of the row with the highest time_updated wins, same as the id");
OpenCodeSessionDiscovery discovery = new OpenCodeSessionDiscovery(root);
assertEquals("gpt-5.6-terra", discovery.actualModelForSessionId("ses_ours").id(),
"querying by id reads OUR row, not the directory's newest row");
assertEquals("deepseek-v4-flash", discovery.actualModelForSessionId("ses_sibling").id(),
"each id resolves to its own row independently of time_updated ordering");
}
@Test
void aNonMatchingDirectoryYieldsUnknownModelRatherThanAMismatch(@TempDir Path root) throws Exception {
void anUnknownSessionIdYieldsUnknownModelRatherThanAMismatch(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"id\":\"x\",\"providerID\":\"y\"}");
assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/other"),
"no row for this cwd yet → unknown, not a wrong model");
assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_no_such_row"),
"no row for this id yet → unknown, not a wrong model");
}
@Test
void aBlankOrNullSessionIdYieldsUnknownModel(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"id\":\"x\",\"providerID\":\"y\"}");
OpenCodeSessionDiscovery discovery = new OpenCodeSessionDiscovery(root);
assertNull(discovery.actualModelForSessionId(null));
assertNull(discovery.actualModelForSessionId(" "));
}
@Test
void aNullModelColumnYieldsUnknownWithoutThrowing(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L, null);
assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"),
assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa"),
"a row with no model value yet is unknown, not a mismatch");
}
@@ -190,7 +201,7 @@ class OpenCodeSessionDiscoveryTest {
void unparseableModelJsonYieldsUnknownWithoutThrowing(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L, "this is not json");
assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"),
assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa"),
"JSON that fails to parse resolves to unknown, never an exception");
}
@@ -198,13 +209,13 @@ class OpenCodeSessionDiscoveryTest {
void modelJsonMissingIdYieldsUnknown(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"providerID\":\"openai\"}");
assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"),
assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa"),
"no id in the JSON → unknown, since id is what a caller actually compares");
}
@Test
void aMissingDatabaseYieldsUnknownModelWithoutThrowing(@TempDir Path root) {
assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"));
assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa"));
}
/**
@@ -212,7 +223,7 @@ class OpenCodeSessionDiscoveryTest {
* shape a real opencode upgrade/downgrade could produce. This must degrade to UNKNOWN for the
* model, and — the property that actually matters — must NOT take id resolution down with it.
* A combined single query for both columns would fail this test; that is why
* {@link OpenCodeSessionDiscovery#actualModelForDirectory} runs its own independent query.
* {@link OpenCodeSessionDiscovery#actualModelForSessionId} runs its own independent query.
*/
@Test
void aMissingModelColumnYieldsUnknownButIdResolutionStillWorks(@TempDir Path root) throws Exception {
@@ -234,7 +245,7 @@ class OpenCodeSessionDiscoveryTest {
OpenCodeSessionDiscovery discovery = new OpenCodeSessionDiscovery(root);
assertEquals("ses_aaa", discovery.sessionIdForDirectory("/w/a"),
"id resolution must survive a database with no model column at all");
assertNull(discovery.actualModelForDirectory("/w/a"),
assertNull(discovery.actualModelForSessionId("ses_aaa"),
"no model column → unknown, not a throw and not a mismatch");
}
}
@@ -56,7 +56,11 @@ class MessageServiceTest {
/** Run {@code send} on a background thread; the current thread drives the worker's turn. */
private CompletableFuture<MessageService.Reply> sendAsync() {
return CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 5000));
return sendAsync("do the task");
}
private CompletableFuture<MessageService.Reply> sendAsync(String content) {
return CompletableFuture.supplyAsync(() -> messages.send(T, content, 5000));
}
private void awaitWaiting() throws InterruptedException {
@@ -86,6 +90,94 @@ class MessageServiceTest {
assertTrue(reply.completed(), "a scraped completion still counts as completed");
}
@Test
void completionFallbackReplacesAnEchoedInjectedBriefWithNoReportOutcome() throws Exception {
String brief = "Implement the requested change. ".repeat(20);
CompletableFuture<MessageService.Reply> send = sendAsync(brief);
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
herdr.readText("⏺ " + brief + "\n❯ ");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome());
assertEquals(CompletionResolver.NO_REPORT_PREFIX + "]", reply.text(),
"the real injector -> completion fallback path must not return the lead's brief");
}
@Test
void completionFallbackKeepsARealReportThatRestatesTheWholeBrief() throws Exception {
String brief = "Load the implementer skill. You own fleetd #999. ".repeat(12);
String report = brief + "\n\n## Report\n"
+ ("I implemented the fix in GitWorktrees.java, added five tests, ran mvn clean install "
+ "and got 1168 tests with 0 failures. Commit 321d8dc pushed. ").repeat(8);
CompletableFuture<MessageService.Reply> send = sendAsync(brief);
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
herdr.readText("⏺ " + report + "\n❯ ");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome());
assertEquals(report.strip(), reply.text(), "a report that restates the whole brief must survive unchanged");
}
@Test
void completionFallbackNamesTheKnownWorktreeAndBranchForAnEchoedBrief() throws Exception {
FakeHerdr localHerdr = new FakeHerdr();
Rendezvous localRendezvous = new Rendezvous();
java.util.concurrent.atomic.AtomicLong localClock = new java.util.concurrent.atomic.AtomicLong();
CompletionResolver localCompletion = new CompletionResolver(new AgentControl(localHerdr), localRendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(),
dev.ltms.fleet.inject.BackendErrorPatternLookup.legacy(),
dev.ltms.fleet.inject.BackendErrorSink.none(),
() -> localClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1),
target -> new CompletionResolver.WorktreeBranch("/tmp/member-worktree", "worker/cb241"));
Injector localInjector = new Injector(new AgentControl(localHerdr), localCompletion);
MessageService localMessages = new MessageService(new AgentControl(localHerdr), localInjector,
localRendezvous, new InMemoryReplyInbox());
String brief = "Implement the requested change. ".repeat(20);
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000));
long deadline = System.currentTimeMillis() + 2000;
while (!localRendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(localRendezvous.isWaiting(T));
localHerdr.readText("$ prompt");
localInjector.onStatus(T, AgentStatus.IDLE);
localInjector.onStatus(T, AgentStatus.WORKING);
localHerdr.readText("⏺ " + brief + "\n❯ ");
localInjector.onStatus(T, AgentStatus.IDLE);
assertEquals(CompletionResolver.NO_REPORT_PREFIX + " worktree=/tmp/member-worktree "
+ "branch=worker/cb241]", send.get(5, TimeUnit.SECONDS).text());
}
@Test
void completionFallbackKeepsTheClippedMarkerWhenAnEchoedBriefIsTooLong() throws Exception {
String brief = "a".repeat(4_001); // CompletionResolver's 4,000-character scrape cap
CompletableFuture<MessageService.Reply> send = sendAsync(brief);
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
herdr.readText("⏺ " + brief + "\n❯ ");
injector.onStatus(T, AgentStatus.IDLE);
assertEquals(CompletionResolver.NO_REPORT_PREFIX + "]\n"
+ "[Pane tail clipped: member did not call fleet_reply.]",
send.get(5, TimeUnit.SECONDS).text());
}
@Test
void backendErrorScrapeThroughMessageServiceFailsInsteadOfBecomingReplyText() throws Exception {
// fleetd#164 (part 2 addendum): a scrape that reads cleanly but is only the backend's own
@@ -43,6 +43,7 @@ class ReplyPushLoopTest {
private static final String PRIMARY = "term_primary";
private static final String WORKER = "term_worker";
private static final String WORKER2 = "term_worker2";
private static final String OTHER_PRIMARY = "term_other_primary";
private static final ObjectMapper MAPPER = new ObjectMapper();
private PrimaryRegistry registry;
@@ -529,6 +530,106 @@ class ReplyPushLoopTest {
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
}
// --- CB-201: backend credential outage nudges -----------------------------------------------
@Test
void oneBackendIncidentForTwoWorkersOfOneLeadSendsOnceAndIsOneShot() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
registry.recordDelegation(WORKER, PRIMARY);
registry.recordDelegation(WORKER2, PRIMARY);
var loop = loop(2, 100);
loop.onBackendIncident("incident-1", List.of(WORKER, WORKER2), "cred-a",
List.of("terra", "codex"), 47);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one lead gets one outage injection");
Thread.sleep(200);
assertEquals(1, rec.sendCount(), "two workers for one lead must send only one outage notice");
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("cred-a"));
assertTrue(nudge.contains("codex, terra"));
assertTrue(nudge.contains(WORKER) && nudge.contains(WORKER2));
assertTrue(nudge.contains("47 remaining seconds"));
assertTrue(nudge.contains("fleet_list"));
loop.onBackendIncident("incident-1", List.of(WORKER, WORKER2), "cred-a",
List.of("terra", "codex"), 47);
Thread.sleep(250);
assertEquals(1, rec.sendCount(), "a delivered incident id must never be sent again");
}
@Test
void oneBackendIncidentSendsOnceToEachAffectedLead() throws Exception {
var rec = recordingClient();
rec.sendLatch = new CountDownLatch(2);
agents = new AgentControl(rec);
registry.recordDelegation(WORKER, PRIMARY);
registry.recordDelegation(WORKER2, OTHER_PRIMARY);
var loop = loop(1, 50);
loop.onBackendIncident("incident-2", List.of(WORKER, WORKER2), "cred-a", List.of("terra"), 60);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "each affected lead gets its one notice");
assertEquals(2, rec.sendCount());
}
@Test
void backendIncidentWaitsForBusyLeadThenCombinesWithFailedTicket() throws Exception {
var rec = new BusyThenIdleHerdrClient(2);
agents = new AgentControl(rec);
var loop = loop(1, 50);
loop.onTicketTerminal("task-failed", WORKER, true);
loop.onBackendIncident("incident-3", List.of(WORKER), "cred-b", List.of("terra"), 31);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "busy lead should receive deferred combined nudge");
assertEquals(1, rec.sendCount(), "ticket and outage must use the same injection");
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("task-failed") && nudge.contains("cred-b"),
"the one nudge must name both the failed ticket and outage: " + nudge);
}
@Test
void failedBackendIncidentSendRetriesWithinTheExistingReminderCap() throws Exception {
var rec = new ThrowOnceHerdrClient();
agents = new AgentControl(rec);
var loop = loop(2, 50);
loop.onBackendIncident("incident-4", List.of(WORKER), "cred-c", List.of("terra"), 20);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "failed first send must retry under maxReminders");
assertEquals(2, rec.promptAttempts.get(), "one failed send and one retry use the existing cap");
}
@Test
void unmappedBackendTargetUsesTheKnownLeadSchedule() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
var loop = loop(1, 50);
loop.onBackendTargetUnmapped(WORKER, "no credential matched");
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = ((Map<?, ?>) rec.sentParams().getFirst().getValue()).get("text").toString();
assertEquals("A backend error on worker term_worker could not be mapped to a credential "
+ "(no credential matched) — no cool-off was applied. "
+ "Run fleet_list to check that worker.",
nudge);
}
@Test
void stopClearsPendingBackendIncidents() {
agents = agentWithStatus("idle");
var loop = loop(1, 100_000);
loop.onBackendIncident("incident-stop", List.of(WORKER), "cred-stop", List.of("terra"), 60);
loop.stop();
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0, 0),
"stop must clear a pending backend incident as well as older sources");
}
// --- CB-590 follow-up: per-source reminder budgets — the regression this round exists for ---
@Test
@@ -928,6 +1029,31 @@ class ReplyPushLoopTest {
return new RecordingHerdrClient();
}
/** Fails its first prompt call, then records the retry. */
private static final class ThrowOnceHerdrClient implements HerdrClient {
private final AtomicInteger promptAttempts = new AtomicInteger();
private final CountDownLatch sendLatch = new CountDownLatch(1);
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
.put("terminal_id", PRIMARY).put("agent_status", "idle"));
}
if ("agent.prompt".equals(method) && promptAttempts.getAndIncrement() == 0) {
throw new IllegalStateException("first prompt fails");
}
if ("agent.prompt".equals(method)) {
sendLatch.countDown();
}
return MAPPER.createObjectNode();
}
@Override
public void close() {
}
}
/**
* Thread-safe fake that reports {@code working} (not injectable) for its first
* {@code busyChecks} status calls, then {@code idle} forever after — used to prove work queued
@@ -0,0 +1,266 @@
package dev.ltms.fleet.placement;
import org.junit.jupiter.api.Test;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #201 / #227: the credential-keyed outage decision itself, isolated from the classifier
* that produces events (Unit 1) and from the lead notice / spawn wiring (Units 3 and 5). The clock
* is a plain {@link AtomicLong} of nanos so the window and cool-off are exercised without a real
* sleep, exactly like {@code BackendQuarantineTest}.
*/
class BackendOutagePolicyTest {
private static final long SECOND = 1_000_000_000L;
@Test
void oneErrorCreatesNoIncidentAndNoCoolOff() {
BackendOutagePolicy policy = new BackendOutagePolicy(() -> 0L);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-1", "API Error: 429");
assertTrue(incident.isEmpty());
assertTrue(policy.remainingCoolOffSeconds("shared-openai").isEmpty());
}
@Test
void twoErrorsWithinTheWindowCreateExactlyOneIncidentAndA60SecondCoolOff() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429").isEmpty());
now.set(30 * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-2", "API Error: 500");
assertTrue(incident.isPresent());
BackendOutagePolicy.Incident value = incident.get();
assertEquals("shared-openai", value.credentialId());
assertEquals(Set.of("session-1", "session-2"), value.targets());
assertEquals(2, value.evidenceCount());
assertEquals(60L, value.windowSeconds());
assertEquals(60L, value.remainingCoolOffSeconds());
assertEquals(OptionalLongOf60(), policy.remainingCoolOffSeconds("shared-openai"));
}
@Test
void repeatedErrorsFromTheSameTargetNeverCreateAnIncidentHoweverManyTimesTheyRepeat() {
// The classifier is a heuristic and a valid member report can quote an "API Error:" line, so
// one target repeating that line inside the window must never, by itself, cost the credential
// its capacity — only a SECOND, DISTINCT target crossing the threshold does.
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
for (int i = 0; i < 5; i++) {
now.set(i * 10L * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-1", "API Error: 429");
assertTrue(incident.isEmpty(),
"the same target repeating must never create an incident on its own (repeat #" + i + ")");
}
assertTrue(policy.remainingCoolOffSeconds("shared-openai").isEmpty());
}
@Test
void evidenceCountIsDistinctTargetsWhileReasonsKeepsEveryEvent() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429").isEmpty());
now.set(10 * SECOND);
// a repeat from the already-counted target: still no incident, still only 1 distinct target
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429 (again)").isEmpty());
now.set(20 * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-2", "API Error: 500");
assertTrue(incident.isPresent());
BackendOutagePolicy.Incident value = incident.get();
assertEquals(Set.of("session-1", "session-2"), value.targets());
assertEquals(2, value.evidenceCount(), "evidenceCount is the distinct-target count");
assertEquals(3, value.reasons().size(),
"reasons keeps every event, including the same-target repeat that did not advance the count");
assertEquals(java.util.List.of("API Error: 429", "API Error: 429 (again)", "API Error: 500"),
value.reasons());
}
@Test
void twoErrorsMoreThanAMinuteApartCreateNoIncident() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429").isEmpty());
now.set(60 * SECOND + 1);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-2", "API Error: 500");
assertTrue(incident.isEmpty(), "the first error's evidence must not survive past the window");
assertTrue(policy.remainingCoolOffSeconds("shared-openai").isEmpty());
}
@Test
void exactly60SecondsApartIsInsideTheWindow() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429").isEmpty());
now.set(60 * SECOND); // exactly the boundary
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-2", "API Error: 500");
assertTrue(incident.isPresent(), "60 seconds apart is inclusive of the window boundary");
}
@Test
void differentCredentialsNeverShareEvidence() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("credential-a", "session-1", "API Error: 429").isEmpty());
now.set(10 * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("credential-b", "session-2", "API Error: 429");
assertTrue(incident.isEmpty(), "an unrelated credential must not be swept into the count");
}
@Test
void twoDifferentProfilesResolvingToTheSameCredentialShareEvidence() {
// The caller is responsible for resolving a profile to its credential id (see
// FleetConfig.Profile#effectiveCredentialId(), same as BackendQuarantine); this class only
// ever sees the resolved credentialId, so two different profile-originated events landing on
// the same credentialId must correlate.
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-from-profile-a", "API Error: 429").isEmpty());
now.set(5 * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-from-profile-b", "API Error: 500");
assertTrue(incident.isPresent());
}
@Test
void errorsDuringCoolOffNeitherExtendTheDeadlineNorReturnAnotherIncident() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
policy.record("shared-openai", "session-1", "API Error: 429");
now.set(1 * SECOND);
policy.record("shared-openai", "session-2", "API Error: 500"); // crosses threshold, cool-off starts at t=1s, until t=61s
now.set(50 * SECOND);
Optional<BackendOutagePolicy.Incident> duringCoolOff =
policy.record("shared-openai", "session-3", "API Error: 500");
assertTrue(duringCoolOff.isEmpty(), "no second incident while cooling off");
assertEquals(OptionalLongOf(11), policy.remainingCoolOffSeconds("shared-openai"),
"the deadline must not have been pushed out by the error during cool-off");
}
@Test
void expiryClearsOldEvidenceSoTwoFreshErrorsAreRequiredToRearm() {
// Both errors land at the SAME instant (t=0) so the window boundary (60s from t=0) and the
// cool-off boundary (60s from the threshold crossing, also t=0) coincide exactly at t=60s.
// That is deliberate: with WINDOW_NANOS == COOLOFF_NANOS, checking "one fresh error" at any
// later instant would also be caught by the window naturally going stale, which would pass
// even if the explicit evidence-clear on cool-off expiry were missing. Checking exactly AT
// t=60s is the one instant where "now - firstEventNanos > WINDOW_NANOS" is false (60 is not
// > 60) but the cool-off has already elapsed ("now < coolOffUntilNanos" is false at 60 >= 60)
// — so only an explicit clear, not a side effect of the window check, keeps this from
// instantly re-mounting the old evidence and minting a bogus incident off one fresh error.
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
policy.record("shared-openai", "session-1", "API Error: 429");
policy.record("shared-openai", "session-2", "API Error: 500"); // cool-off until t=60s
now.set(60 * SECOND); // cool-off just elapsed, exactly at the window's own boundary too
Optional<BackendOutagePolicy.Incident> oneFreshError =
policy.record("shared-openai", "session-4", "API Error: 500");
assertTrue(oneFreshError.isEmpty(), "a single fresh error must not immediately re-trip");
now.set(61 * SECOND);
Optional<BackendOutagePolicy.Incident> secondFreshError =
policy.record("shared-openai", "session-5", "API Error: 500");
assertTrue(secondFreshError.isPresent(), "two fresh errors after expiry re-arm the policy");
}
@Test
void remainingSecondsRoundUp() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
policy.record("shared-openai", "session-1", "API Error: 429");
now.set(1); // 1 nanosecond later, still crosses the threshold
policy.record("shared-openai", "session-2", "API Error: 500");
// cool-off deadline is now(=1ns) + 60s exactly; a fraction-of-a-second-remaining check:
now.set(1 + 59 * SECOND + 1); // 59.000000001s of the cool-off has passed -> < 1s remains
assertEquals(OptionalLongOf(1), policy.remainingCoolOffSeconds("shared-openai"),
"a sub-second remainder must round up to a full second, matching BackendQuarantine");
}
@Test
void aConcurrentSecondAndThirdCallReturnOnlyOneIncidentBetweenThem() throws InterruptedException {
for (int attempt = 0; attempt < 25; attempt++) {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
int threads = 8;
ExecutorService pool = Executors.newFixedThreadPool(threads);
CountDownLatch ready = new CountDownLatch(threads);
CountDownLatch go = new CountDownLatch(1);
AtomicInteger incidents = new AtomicInteger();
try {
for (int i = 0; i < threads; i++) {
int idx = i;
pool.submit(() -> {
ready.countDown();
try {
go.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-" + idx, "API Error: 500");
if (incident.isPresent()) {
incidents.incrementAndGet();
}
});
}
assertTrue(ready.await(5, TimeUnit.SECONDS));
go.countDown();
pool.shutdown();
assertTrue(pool.awaitTermination(5, TimeUnit.SECONDS));
assertEquals(1, incidents.get(),
"exactly one incident must be minted across every concurrent caller on attempt " + attempt);
} finally {
pool.shutdownNow();
}
}
}
private static java.util.OptionalLong OptionalLongOf60() {
return OptionalLongOf(60);
}
private static java.util.OptionalLong OptionalLongOf(long value) {
return java.util.OptionalLong.of(value);
}
}
@@ -300,6 +300,151 @@ class SessionManagerTest {
assertTrue(sessions.roster().contains(updated), "FAILED is still in acquired-minus-released roster");
}
@Test
void backendErrorWinsBeforeNormalCompletionFromBusy() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
sessions.onTurnComplete(session.terminalId());
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state(),
"normal completion must not restore DONE after a backend error from BUSY");
assertEquals("backend exited", updated.failureReason());
}
@Test
void backendErrorWinsAfterNormalCompletionFromDone() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
sessions.onTurnComplete(session.terminalId());
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state(),
"backend error must replace DONE when normal completion won first");
assertEquals("backend exited", updated.failureReason());
}
@Test
void backendErrorMemberCannotReceiveAnotherDelivery() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
assertEquals(MemberSession.State.BACKEND_ERROR, sessions.get(session.paneId()).orElseThrow().state(),
"onDelivered must refuse a backend-error member");
}
@Test
void losingCompletionDoesNotReleaseOrClearABackendErrorMember() {
assertLosingCompletionLeavesBackendErrorIntact(1, false,
"a stale completion must not release a backend-error member at the context cap");
assertLosingCompletionLeavesBackendErrorIntact(0, true,
"a stale completion must not clear the context of a backend-error member");
}
private void assertLosingCompletionLeavesBackendErrorIntact(int contextCap, boolean clearAfterTurn,
String protectedAction) {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
ClearContextSpyLauncher launcher = new ClearContextSpyLauncher(delegate);
java.util.concurrent.atomic.AtomicReference<SessionManager> manager = new java.util.concurrent.atomic.AtomicReference<>();
java.util.concurrent.atomic.AtomicReference<String> terminal = new java.util.concurrent.atomic.AtomicReference<>();
java.util.concurrent.atomic.AtomicBoolean armed = new java.util.concurrent.atomic.AtomicBoolean();
SessionManager sessions = new SessionManager(launcher, new RecordingWorktrees(), () -> {
if (armed.compareAndSet(true, false)) {
manager.get().onBackendError(terminal.get(), "backend exited during completion");
}
return 1;
}, contextCap, clearAfterTurn);
manager.set(sessions);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
terminal.set(session.terminalId());
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
armed.set(true);
assertFalse(sessions.onTurnCompleteWithPostAction(session.terminalId()),
"the completion CAS loses after the injected backend error");
assertTrue(sessions.get(session.paneId()).isPresent(), protectedAction);
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state());
assertEquals("backend exited during completion", updated.failureReason());
assertFalse(herdr.called("pane.close"), protectedAction);
assertEquals(0, launcher.clearContextCalls(), protectedAction);
}
@Test
void backendErrorForUnknownTargetIsWarnedAndDoesNotCreateASession() {
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger sessionLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
sessionLog.addAppender(appender);
try {
SessionManager sessions = sessionManager(new FakeHerdr());
assertFalse(sessions.onBackendError("term_missing", "backend exited"));
assertTrue(sessions.roster().isEmpty(), "unknown target must not create a session");
assertTrue(appender.list.stream().anyMatch(e -> e.getLevel().equals(Level.WARN)
&& e.getFormattedMessage().contains("term_missing")),
"unknown target is logged at WARN");
} finally {
sessionLog.detachAppender(appender);
}
}
@Test
void backendErrorForReleasedTargetDoesNotRecreateTheSession() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.release(session.paneId());
assertFalse(sessions.onBackendError(session.terminalId(), "backend exited"));
assertTrue(sessions.get(session.paneId()).isEmpty(), "released session must stay absent");
}
@Test
void rosterViewIncludesFailureReasonOnlyForBackendError() {
MemberSession failed = new MemberSession("pane-error", "term-error", "prof", MemberRole.DEV,
"/cwd", null, 0, 0, 1, MemberSession.State.BACKEND_ERROR, null, null,
null, null, "backend exited");
MemberSession ordinary = new MemberSession("pane-ready", "term-ready", "prof", MemberRole.DEV,
"/cwd", null, 0, 0, 0, MemberSession.State.READY, null, null);
Map<String, Object> failedView = SessionManager.rosterView(failed, null);
Map<String, Object> ordinaryView = SessionManager.rosterView(ordinary, null);
assertEquals("backend_error", failedView.get("state"));
assertEquals("backend exited", failedView.get("failureReason"));
assertFalse(ordinaryView.containsKey("failureReason"),
"ordinary rows must not gain a blank failureReason field");
}
@Test
void onTurnFailedIsLoggedAtWarnWithThePriorState() {
// CB-564: this transition used to be a bare DEBUG "session marked failed" — a symptom with no
@@ -1004,6 +1149,76 @@ class SessionManagerTest {
}
}
/** Delegates all launcher work while recording context clears. */
private static final class ClearContextSpyLauncher implements PeerLauncher {
private final PeerLauncher delegate;
private int clearContextCalls;
ClearContextSpyLauncher(PeerLauncher delegate) {
this.delegate = delegate;
}
@Override
public Set<Capability> capabilities() {
return delegate.capabilities();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return delegate.capabilitiesFor(profileName);
}
@Override
public PeerHandle spawn(SpawnRequest req) {
return delegate.spawn(req);
}
@Override
public Set<String> profiles() {
return delegate.profiles();
}
@Override
public String defaultProfile() {
return delegate.defaultProfile();
}
@Override
public String effectiveCwd(SpawnRequest req) {
return delegate.effectiveCwd(req);
}
@Override
public List<String> parityOverlay(String profileName) {
return delegate.parityOverlay(profileName);
}
@Override
public List<?> list() {
return delegate.list();
}
@Override
public int reapOrphanWorkers() {
return delegate.reapOrphanWorkers();
}
@Override
public void stop(String id) {
delegate.stop(id);
}
@Override
public boolean clearContext(String id) {
clearContextCalls++;
return delegate.clearContext(id);
}
int clearContextCalls() {
return clearContextCalls;
}
}
@Test
void acquireRecordsTheAgentSessionIdFromTheHandle() {
FakeHerdr herdr = new FakeHerdr();