Compare commits

...

28 Commits

Author SHA1 Message Date
Dai Ha 69e09b10fa #388: scrub a pane shell that is neither login nor interactive
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 1m32s
EnvAllowListScrub generated four zsh startup files but only .zshrc and
.zlogin sourced the scrub — .zshenv (the one file zsh always reads) did
not. A pane shell that is neither login nor interactive reads only
.zshenv and stops, so it was never scrubbed at all (measured on fleet01,
issue #388).

Adding an unguarded scrub to .zshenv (the ticket's own suggested fix) is
wrong: .zshenv is read by every zsh, including a short-lived `zsh -c`
a member's own tooling forks for a single command. Those children are
also neither login nor interactive, so they would scrub the environment
their parent deliberately set for them (GIT_DIR, VIRTUAL_ENV, ...), and
the rewritten scrub-report.txt would describe the last child to exit
instead of the pane.

Fix (per comment 15387, measured): keep .zshrc/.zlogin unconditional,
and add to .zshenv a pass guarded on the exact condition that defines
the gap (neither login nor interactive), plus a per-pane sentinel
(_CB633_SCRUBBED) so it runs once per pane, not once per process. The
sentinel is exported only after the scrub runs, and is folded into the
scrub's own allow-list so a later pass in the same pane cannot blank it
back to empty.

Also corrects the class javadoc's wrong premise (a bare argv[0] proves
NOT login, not "therefore interactive") and its now-stale two-file
walkthrough.

Tests: two new real-zsh tests in EnvAllowListScrubTest run actual
non-login/non-interactive zsh processes (never string-match the
generated files) to prove: a neither-shell pane is scrubbed; a child
that pane forks keeps variables the pane deliberately set for it; the
child does not re-scrub; and scrub-report.txt still describes the pane
after the child exits. Both fail without the production fix (verified
by reverting it and re-running: AssertionFailedError on the sentinel
and on the decoy secret surviving).
2026-09-10 07:21:08 +07:00
Dai Ha b9d09e044e t386: pin the per-member drift baseline the fix's own tests left open
CI / contract (push) Successful in 37s
CI / build (push) Successful in 1m57s
The two tests merged with #386 both start with the member already BUSY, so a
single global drift baseline passes them. This one sleeps the host while nothing
is busy and only then starts a turn, which fails without the per-member map.
2026-09-10 06:59:42 +07:00
Dai Ha fd8650cda4 Merge #386: correct the stall check for a monotonic clock frozen by host sleep 2026-09-10 06:56:37 +07:00
Dai Ha 11050e24ed t384: fix javadoc indentation on the merged shareWithGroup lines
CI / contract (push) Successful in 1m7s
CI / build (push) Successful in 1m29s
2026-09-10 06:54:51 +07:00
Dai Ha 769f282408 #386: give the stall detector a real-time clock, log the divergence
CI / contract (pull_request) Successful in 1m18s
CI / build (pull_request) Successful in 1m26s
FleetHealthMonitor.tick's stalled check compared two monotonic-clock
readings (System.nanoTime(), which macOS freezes across a host sleep),
so a member BUSY for 101 real minutes was never flagged.

The monitor now also takes a wall-clock LongSupplier (realtimeClock),
used only inside the stall check. Each tick measures how far the two
clocks moved apart since the previous tick and folds any positive
divergence into a running total; when a single tick's divergence
exceeds one tick interval (the signature of a sleep, since a tick
cannot run while the process itself is suspended) it logs one WARN
naming how long the detector could not see. The correction is applied
per member, keyed to when that member's current lastActivityAtNanos
was first observed BUSY - not since monitor start - so a sleep that
happened before a member went busy is never charged to it.

Every other use of the monitor's clock (readiness grace, snapshot
timestamp) is unchanged. Backend quarantine/cool-off, the lead tab
scan, the completion resolver, the session reaper and the message
service TTLs are untouched, per the ticket's decision.

Existing FleetHealthMonitor/FleetHealth tests pass unmodified (none of
them ticks a BUSY session more than once, so the drift path never
engages for them). Two new tests: a frozen monotonic clock past the
real-time threshold produces STALL_SUSPECTED, and a single sleep gap
logs the divergence exactly once, not once per tick.
2026-09-10 06:53:30 +07:00
Dai Ha 6f71f40047 Merge #384: pre-create the scrub receipt and give it group write 2026-09-10 06:49:46 +07:00
Dai Ha fb36c5238f Merge #382: SpawnRequest.withProfile() replaces the six-accessor rebuild 2026-09-10 06:49:46 +07:00
Dai Ha cc9cdc938b #384: write shared scrub receipts 2026-09-10 06:45:30 +07:00
Dai Ha 5fede82468 #382: give SpawnRequest a withProfile wither, guard it against the arity trap
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m49s
CompositePeerLauncher:372 rebuilt a routed SpawnRequest from a literal
new SpawnRequest(...) call listing six of the original request's own
accessors. That call is only correct because it happens to match the
canonical 6-arg constructor today; add a 7th component plus the
established back-compat constructor at the old (now-shorter) arity and
this call would silently rebind to it, dropping the new field on every
profile-routed spawn with no compile error — the same defect shape
already guarded on FleetConfig.withDefaults() (#357) and MemberSession
(#358).

Add SpawnRequest.withProfile(String), modeled on
MemberSession.withState/withActivity, and use it at the call site
instead. Add a guard test that resolves the true canonical constructor
by exact component types (never by argument count), gives every
component a distinctive value, and asserts every component but
profileName survives withProfile() unchanged.

Proved the guard against the real mechanism: temporarily dropped the
last (role) argument from withProfile()'s constructor call so it bound
to the 5-arg back-compat constructor — it still compiled, and the new
test failed, catching the silently-defaulted role. Restored the fix
and reconfirmed green.
2026-09-10 06:44:15 +07:00
Dai Ha 1515025804 charter: a profile differs in liveness, not just model and cost
CI / contract (push) Successful in 1m12s
CI / build (push) Successful in 1m31s
The old line said profiles differ in model and cost. That is true and it is
not the reason the default hurts. Measured on two hosts: a default sitting on
an exhausted or withdrawn credential either fails the spawn loudly or, worse,
produces a member that starts fine and then returns nothing.

Wording proposed by the fleet01 lead; merged with the existing 'not in tier'
clause, which is still right. Applied byte-identically to the wiki template.
2026-09-10 06:28:51 +07:00
Dai Ha 7754f53662 Merge t385 follow-up: pin the mid-turn redelivery ack 2026-09-10 06:28:08 +07:00
Dai Ha a507f7b31b t385: pin that a redelivery is acked even while the lead is mid-turn
Verification found the fix's dedup check could be moved behind the
injectable-status gate and every test still passed. That placement matters:
this lead is mid-turn most of the time, so gating the ack on an idle pane
leaves the redelivered message held, and the next recovery delivers it again.

The new test fails on that mutant and passes on the fix.
2026-09-10 06:28:03 +07:00
Dai Ha 71c322f104 Merge #385: a redelivered lead message must not write the pane twice 2026-09-10 06:24:40 +07:00
Dai Ha fde2c15627 t385: a redelivered lead message must not write the pane twice
An AMQP recovery clears the held delivery tags, the broker redelivers with
fresh ones, and the coordination loop wrote the same peer message into the
lead pane again. Measured on the live daemon: one msgId reached the pane 12
times in 9 hours, across 19 recovery events.

LeadCoordLoop now remembers the msgIds it has written to a pane (bounded at
1024) and acks a redelivery without a second write. LeadMailbox.ack no longer
returns quietly for an unknown msgId: it throws, so an ack that never reached
the broker is reported instead of hidden. A repeat ack that this connection
already completed stays quiet, tracked in a bounded set.
2026-09-10 06:24:35 +07:00
Dai Ha 2830735644 t377: raise the logger level in the test, or it asserts on an empty list
CI / contract (push) Successful in 46s
CI / build (push) Failing after 1m32s
The 9 tests the member wrote failed 7 of 9 on first build. The production
code was correct; the test harness was not. logback-test.xml sets
dev.ltms.fleet to WARN, so the INFO shape lines were dropped by the level
check before any appender saw them.

The member copied attach()/detach() from MemberTrustModelReportTest but not
the setLevel(INFO) those siblings do at each call site. Doing it inside
attach()/detach() covers all nine at once, and restores the original level
(null, meaning inherit) rather than a concrete one.

Proven by mutation: making the code log the value fails 4 tests, including
theValueNeverAppearsInLogOutput — 'the full GITEA_HOST value must never
reach the log'. Full build 1459 tests green.
2026-09-09 07:52:43 +07:00
Dai Ha 4e3ac91a22 Merge worker/t377-7f587b-7: startup reports the git host value's shape, never its value 2026-09-09 07:49:09 +07:00
Dai Ha 29d3f0b41f t377: report the git host value's SHAPE at startup, never its value
A member receives the git host as GITEA_HOST and often gets a full URL
(scheme, trailing slash) where it expects a bare host, so it builds
https://https://... and the request never leaves the machine.

Logs set/unset, length, startsWithScheme and trailingSlash next to the
existing secret report. The value itself is never logged, and the value
is still passed to members unchanged — rewriting it here would change
what works on one host and breaks on another.

Written by the gx member on branch worker/t377-7f587b-7; it could not
build, commit or push because the host command classifier refused every
shell command (see #381). Build and verification are mine.
2026-09-09 07:49:02 +07:00
Dai Ha cbc732444f fleetd #365: the push-nudge metric outcome is 'sent', not 'delivered'
CI / contract (push) Successful in 1m22s
CI / build (push) Successful in 1m31s
The label was renamed in the #365 merge. This table still named the old one.
'sent' counts the herdr paste-and-submit call returning, never a confirmation
that the pane read it.
2026-09-09 07:40:57 +07:00
Dai Ha 6f828b8c38 Merge worker/t365-3920c5-3: t365: fleet_reply/REST reply distinguish resolved vs queued; rename nudge metric outcome
CI / contract (push) Successful in 46s
CI / build (push) Successful in 1m48s
2026-09-09 07:35:00 +07:00
Dai Ha 145a8c8862 Merge worker/t373-336973-2: t373: pin the production XDG-excludes seam GitWorktreesTest.seedingGitWorktrees builds 2026-09-09 07:35:00 +07:00
Dai Ha 7057291739 Merge worker/t358-6e989b-1: t358: guard MemberSession's 5 rebuild sites and Profile.withProfile() against the back-compat-arity trap 2026-09-09 07:35:00 +07:00
Dai Ha 2302b3bc11 #376: the too-fast failure stops asserting a cause it cannot know
CI / contract (push) Successful in 50s
CI / build (push) Successful in 1m31s
A turn that settles inside MIN_TURN_NANOS is still reported FAILED. Only the claim
about WHY is withdrawn, and the pane is still carried.

Measured on fleet01 on 2026-09-08 UTC: an opencode member on mimo-v2.5-free answered
a real question in 1575ms, below the 2000ms floor. The daemon reported the turn as
failed with 'most likely a backend error before any work started'. The answer was
right there in the scrape. With no errorPattern configured — the live state on both
hosts, which both log as 'backend-error classification: off' — that sentence is a
guess, and a reader who believes it stops looking at the pane.

WHAT I REJECTED, because the next person will try it. A worker implemented the
ticket's first suggested direction: inspect the pane inside the floor and resolve a
COMPLETION when the text looks like a real reply. Its test for 'looks like a real
reply' was non-blank plus a '.', '!' or '?' anywhere in the text. That is unsafe
twice over. lastAssistantBlock falls back to the WHOLE pane when it finds no U+23FA
marker, so on a crash the candidate reply is the entire screen; and a crash pane
almost always contains a full stop, in a file path, a version or a hostname. I ran
that implementation against the new guard test and it resolved

    Error: connection reset while loading src/main/java/Foo.java v1.2.3

as a COMPLETION — expected: <FAILED> but was: <COMPLETION>. A loud wrong answer
became a silent one, which is the trade the ticket brief forbade.

The obvious repair does not work either. Requiring the U+23FA marker as positive
evidence would be safe, but that marker is Claude Code chrome and an opencode pane
never carries it — and an opencode member is what raised this ticket. There is no
reliable cross-backend marker for 'this is a real reply', so this path must not try
to judge one. That reasoning is now in the failTooFast javadoc.

Two guard tests. aPlausibleLookingReplyInsideTheFloorStillFails pins the safety
property against exactly the rejected approach; it passes today and fails against
that implementation, which is how it was verified rather than assumed.
theTooFastFailureDoesNotAssertACauseItCannotKnow pins the wording.

One existing test changed. aNonMatchInsideTheFloorStaysGenericAndNeverNotifiesTheSink
asserted the phrase 'too fast to be real work', which carried the withdrawn claim. It
now asserts what it was really guarding: the floor alone fails the turn, the reason
stays generic, the pane is carried, and the typed sink is never notified.

MIN_TURN_NANOS is unchanged at 2000ms. mvn clean install: 1441 tests, 0 failures.
2026-09-09 07:28:04 +07:00
Dai Ha 3a004dc1b3 t373: pin the production XDG-excludes seam GitWorktreesTest.seedingGitWorktrees builds
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m32s
fleetd #362 review finding 2 protects GitWorktrees#previouslyEffectiveExcludesFileContent's
Java-side XDG_CONFIG_HOME/HOME read (it never goes through a git subprocess, so no
GIT_CONFIG_GLOBAL/GIT_CONFIG_SYSTEM isolation reaches it) with a gitEnv constructor seam. A
mutation run during the #372/#369 merge found that seam unpinned: stripping hermeticGitEnv(tmp)
from seedingGitWorktrees left every test green, poisoned XDG_CONFIG_HOME or not.

Adds seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory, which asserts
the property directly (a GitWorktrees built for seeding resolves the fallback inside its own
throwaway directory) using a self-contained marker instead of relying on an externally poisoned
env var. Refactors hermeticGitEnv/seedingGitWorktrees into two-argument overloads (one taking an
explicit XDG_CONFIG_HOME / gitEnv) so the new test can pre-populate the marker before construction
while still going through the same production construction every other seeding test uses; no
behavior change for the 4 existing call sites.
2026-09-09 07:24:22 +07:00
Dai Ha a9a3c12232 t365: fleet_reply/REST reply distinguish resolved vs queued; rename nudge metric outcome
CI / contract (pull_request) Successful in 1m17s
CI / build (pull_request) Successful in 1m51s
fleetd #365. fleet_reply always returned the literal "delivered" and
POST /sessions/{id}/reply always returned {"delivered": true}, whether
the reply resolved a live waiting send/ticket or was merely queued in
the inbox for a later drain (CB-307) — both are successes, but not the
same fact.

MessageService.reply() now returns a ReplyOutcome (RESOLVED_SEND,
RESOLVED_ASYNC_TICKET, or QUEUED) instead of an always-true boolean.
FleetMcp.reply and FleetApp.replyMessage both read it: the MCP tool
result names which happened, and the REST body's "delivered" field is
now accurate, with an added "outcome" field.

Also renames the heartbeat/push-loop nudge metric's "delivered" outcome
to "sent" (LeadHeartbeatLoop, ReplyPushLoop, FleetMetrics): it only
records that the herdr agent.prompt paste-and-submit call succeeded,
never that the lead's pane actually read it — there is no read-receipt
concept at that layer, so "delivered" overclaimed there too.

Tests: MessageServiceTest/FleetMcpTest/FleetAppTest strengthened to
assert the specific outcome per case (a resolved send, a resolved async
ticket, and a queued reply); ReplyPushLoopTest updated for the outcome
rename.
2026-09-09 07:23:41 +07:00
Dai Ha 94ec77a1bc t358: guard MemberSession's 5 rebuild sites and Profile.withProfile() against the back-compat-arity trap
CI / contract (pull_request) Successful in 1m23s
CI / build (pull_request) Successful in 1m53s
Follow-up to #357 (FleetConfig.withDefaults()). Reflectively enumerate each
record's own components, resolve the canonical constructor by exact
component types, build a real non-null value per component, run each
rebuild site, and assert every component survives (except the one it is
documented to change). Exclusion lists pinned at 0 for both.

Re-counted the Profile back-compat ladder directly against the source:
8 constructors (arities 25, 24, 22, 20, 18, 15, 14, 12) against a
canonical arity of 26 — the ticket's own number was explicitly untrusted.
2026-09-09 07:19:12 +07:00
Dai Ha 49f285cfda charter: a measured fact in an addendum must carry its own deletion trigger
CI / contract (push) Successful in 49s
CI / build (push) Successful in 1m36s
The operator chose this and its size; the reasoning below is the fleet01 lead's.

Placed in the canonical block's boundary paragraph rather than the orchestration
body. That paragraph already talks about the addendum layer instead of protocol,
every project that mounts the bridge inherits it, and it sits about 3800 characters
before the primary's step list, so it does not dilute the steps a lead reads while
working. Perishability is structurally an addendum problem: the block is
byte-identical across projects by construction, so a dated local measurement in the
block body would already be a layering violation.

What happened. The fleet01 lead's kb addendum held a dated merge-refusal section
that carried an instruction to delete itself once it stopped reproducing. On
2026-09-08 UTC the operator granted merge rights on akb/kb, the lead re-ran the
probe, got 409 'head out of date' where the identical request had returned 405
'User not allowed to merge PR' on 2026-09-06, and deleted the section as instructed.

Why four parts and not one. The lead's finding is that the banner did not work
because it was emphatic. It worked because the falsification condition was
executable: it carried the exact probe, the reason for the all-zeroes
head_commit_id, and what each response code meant. The lead did not have to
reconstruct the experiment or decide what would count as refutation, and just ran
it. A banner saying 'this may be out of date, verify before relying on it' costs
the same space and does nothing, because deciding what would falsify a claim is the
expensive step and a reader in the middle of another task will not pay it. So: the
date, the command, what each outcome means, and the instruction to delete. The
fourth without the second is decoration.

The closing clause is the justification for the machinery. Most stale notes are
merely wrong. This one went stale in the dangerous direction: it would have told a
future lead it could not merge at the exact moment merging became its job, silently
and with confidence. A note that goes harmlessly stale does not need this.

Note what is NOT centralized here. The banner text itself cannot be. What fired for
the lead was a specific instruction sitting on top of the specific stale fact, which
it could not read past on its way to acting. A rule elsewhere saying 'date your
measurements' would not have fired, because nobody reads that rule at the moment
they re-measure. This sentence sets the convention; the trigger still has to live
next to the fact it governs.

Propagated to the wiki template in the same turn, wiki 8c4f152 on main; the sync
check in this file's addendum reports 'in sync: True'.
2026-09-09 04:24:17 +07:00
Dai Ha 5a12ae7930 charter: test a refusal, and do not count a transport failure as one
CI / contract (push) Successful in 1m11s
CI / build (push) Successful in 1m38s
Step 8 gained a refusal paragraph in 2f71a30, which said what a lead does once the
forge refuses a merge. It did not say how a lead establishes that it was refused.
Both halves of this amendment come from the fleet01 lead, measured on akb/kb on
2026-09-08 UTC, and both are ways to be wrong about a permission you never tested.

Do not read a refusal off a permissions field. After the operator granted merge
rights, the lead re-ran its probe: POST .../pulls/53/merge with an all-zeroes
head_commit_id, chosen so the request cannot succeed on its merits and a rejection
can only mean the refusal. It returned 409 'head out of date' where the identical
request returned 405 'User not allowed to merge PR' on 2026-09-06. A 409 is payload
validation and sits after the permission gate, so the grant took. The lead reports
the repository permissions object did not change across that flip -- still
admin:false, push:true, pull:true. I did not read that object myself; my forge token
is a different identity and would return a different one, so this stays the lead's
measurement and not mine. Merge rights on a protected branch live in branch
protection, so a permissions field can be wrong in both directions.

Do not count a transport failure as a refusal. The lead's first attempt returned
HTTP 000, because GITEA_HOST already carries a scheme and a trailing slash and the
URL came out as https://https://git.ltms.dev//api/... Under a 'not 200' test that is
indistinguishable from being refused. A probe exists to separate a refusal from
everything else, so an error that never reached the gate has to be a third answer
that concludes nothing.

Propagated to the wiki template in the same turn, wiki 8c2ef96 on main; the sync
check in this file's addendum reports 'in sync: True'.
2026-09-09 04:09:08 +07:00
Dai Ha 127e6832a9 Merge #374: fleetd holds off idle sleep while any member is live
CI / contract (push) Successful in 1m18s
CI / build (push) Successful in 2m3s
Lands PR #355 (fleetd #354's sibling), rebased onto current main by a worker
after 39 commits of drift left it unmergeable.

The problem, measured on the original branch: a fleetd host idle-slept after as
little as one minute (pmset -g custom reported 'sleep 1' on battery). Overnight
the daemon's AMQP link dropped 13 times, and every drop minute had a sleep or
wake event in pmset -g log in the same minute or the one before. The AMQP churn
is the visible symptom; the real cost is a member mid-turn freezing with the
host, and a long turn with nobody typing is exactly the case that goes idle.

IdleSleepGuard holds an OS-level assertion for as long as at least one member is
live. It is driven by SessionManager's existing onAcquire/onRelease hooks rather
than a second member count kept in parallel, so it reads the same registry
fleet_list's numbers come from, and only a real 0->1 or 1->0 crossing touches the
OS. It fails safe: a mechanism that cannot acquire means nothing is ever held,
and it never throws, never blocks a spawn, a release, or shutdown.

Conflict resolution was the whole job, and all three were in config plumbing:
ConfigRef, FleetConfig and ConfigRefTopLevelReportingCoverageTest. The power
package is byte-identical to the original branch commit.

Verified on this merge, not taken from the worker's report:
  mvn clean install -> Tests run: 1439, Failures: 0, Errors: 0, BUILD SUCCESS
  (1425 on main + 14 new: 4 caffeinate, 5 guard, 1 wiring, 4 config)

The denominator recount, which the worker flagged as its own weakest number
because this file's count has drifted three times before (#330/#333/#337). I
counted it mechanically rather than reading it: FleetConfig has 24 canonical
record components; COLD_KEYS 5, DEFERRED_KEYS 13, SPLIT_KEYS 3, plus the 3 the
javadoc names as hot-excluded (placement, memberCredentials, memberLoginShell).
5+13+3+3 = 24. The javadoc's '24 components: 5 cold, 13 deferred, 3 split, 3
hot-excluded' is correct. The worker's prose called idleSleepGuard the 25th
constructor argument; it is the 24th. The code is right, the report was off by
one.

Mutation run on merge, on the half the worker verified by READING rather than by
proving -- it said it had checked that withDefaults()'s final call binds the true
canonical constructor. I dropped the trailing idleSleepGuard argument so the call
silently binds the 23-arg back-compat overload. It compiles, which is the whole
hazard. Caught: 1 failure, 3 errors, BUILD FAILURE, and
FleetConfigWithDefaultsPreservesEveryComponentTest names the dropped component
and prints its own denominator -- '24 components, 24 checked, 0 excluded, 23
survived'. That test was added on main after this exact defect happened live when
idleSleepGuard was added on a sibling branch; the worker had to add the missing
entry to it, and doing so is what makes the guard cover this component at all.
2026-09-07 20:31:41 +07:00
31 changed files with 1832 additions and 111 deletions
+19 -2
View File
@@ -7,6 +7,14 @@
> wiki ([Use Cases](https://git.ltms.dev/fleet/fleetd/wiki/7-Use-Cases) → *The portable
> CLAUDE.md block*); improvements go to the template first, then out to each project. Anything
> specific to *this* repo lives under §Project addendum below, never inline above it.
>
> **Anything you measure in an addendum is perishable.** Date it, give the command that
> re-measures it and what each outcome means, and tell the reader to delete the section once
> it stops reproducing. The four parts work together: deciding what would falsify a claim is
> the expensive step, and a reader in the middle of another task will not pay it, so a bare
> "verify before relying on this" costs the same space and does nothing. The case this is for
> is a note that goes stale as a live restriction — it will tell a future session it cannot do
> the thing at the moment doing it becomes the job.
If no `fleet_*` MCP tools are mounted in this session, this section does not apply — skip it.
@@ -71,8 +79,11 @@ below are the procedure — run them in order, every task, not only the big ones
the final judgment call, verification, merges, and anything that depends on context only you
hold. Nothing else is yours by default.
3. **Spawn every delegated unit first** — `fleet_spawn{profile, worktree:true, ticket}`, one per
unit, *before* sending any. Pass `profile` explicitly: profiles differ in model and cost, not in
tier, so the default is rarely what you want.
unit, *before* sending any. Pass `profile` explicitly: profiles differ in model, cost and
LIVENESS, not in tier, so the default is rarely what you want. The default is whatever the
daemon reports, and on a host where it sits on an exhausted or withdrawn credential every
unqualified spawn fails — sometimes loudly, sometimes as a member that spawns fine and then
produces nothing. `fleet_profiles` reports the default; check it once per session.
4. **Then send them all** — `fleet_send{sessionId, content, wait:false}`. Line 1 of every brief is
`Load the <name> skill.` naming the worker's playbook; those skills are opt-in and that line is
what makes them reliable. Where the project ships no such skill, spell the procedure out in the
@@ -100,6 +111,12 @@ below are the procedure — run them in order, every task, not only the big ones
without having read the diff yourself. A refusal is exactly when that shortcut is tempting,
because no action is left that forces you to look, and taking it turns this step into
forwarding a reviewer's verdict — which is delegating the merge by proxy, two lines above.
**Test a refusal; do not read it off a permissions field.** A protected branch holds its merge
rights separately from the repository permissions, so that field can say yes while the merge is
refused, and still say no after a grant makes it work. Probe instead, with a request that cannot
succeed on its merits, so a rejection can only mean the refusal. Treat a transport failure as a
third answer that proves nothing: a timeout, a DNS error or a bad URL is not a refusal, and
counting it as one makes you sure of something you never measured.
**Steps 3 and 4 are separate on purpose** — spawning and sending in one loop is how parallel work
silently becomes serial, and it is the most common way this layer is wasted. For the same reason,
+1 -1
View File
@@ -192,7 +192,7 @@ Deliberately small; every one maps to a failure mode we have actually hit.
| `fleet_send_duration_seconds` | histogram | delegated turn latency |
| `fleet_replies_total{path}` | counter | path ∈ rendezvous\|inbox — how often a reply strands (CB-307's whole reason to exist) |
| `fleet_inbox_depth{target}` | gauge | undrained replies; steady-state should be 0 |
| `fleet_push_nudges_total{outcome}` | counter | outcome ∈ delivered\|exhausted — a rising `exhausted` means the primary is not draining |
| `fleet_push_nudges_total{outcome}` | counter | outcome ∈ sent\|exhausted — a rising `exhausted` means the primary is not draining. `sent` was called `delivered` until fleetd #365; it counts the herdr paste-and-submit call returning, never a confirmation the pane read it |
| `fleet_spawns_total{kind,outcome}` | counter | outcome ∈ ready\|timeout\|guard_rejected; per peer kind (CB-402) |
| `fleet_sessions{state}` | gauge | SPAWNING/READY/BUSY/DONE census |
| `fleet_herdr_calls_total{method,outcome}` | counter | socket health — the dependency everything rests on |
@@ -123,6 +123,11 @@ public final class Fleetd {
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
// boots fine either way — this is the only thing that says so out loud.
reportRequiredSecrets(cfg);
// fleetd #377: the git host value a member receives as GITEA_HOST is often a full URL
// (scheme and trailing slash), not a host name. A member that assumes a bare host then
// builds https://https://... and the request never leaves the machine. Report the shape
// next to the secret report — shape only, never the value.
reportGitHostShape(cfg);
reportMemberTrustModel(cfg);
// CB-596: an absent (or empty) memberCredentials: block blocks NOTHING — no credential
// name is hardcoded any more to fall back on. Say so loudly, the same way a missing
@@ -579,8 +584,12 @@ public final class Fleetd {
if (cfg.health() != null && cfg.health().isEnabled()) {
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
// through the same idempotent target-wide operation CB-516 already uses on release.
// fleetd #386: System::nanoTime freezes across a macOS sleep, so the stall check also
// gets a wall-clock source to detect and correct for that freeze. Every other decision
// in FleetHealthMonitor stays on the monotonic clock, unchanged.
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, cfg.health().intervalOrDefault(),
System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()),
cfg.health().intervalOrDefault(),
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
@@ -1195,6 +1204,76 @@ public final class Fleetd {
});
}
/**
* fleetd #377: the env var names holding the git host value that members receive as
* {@code GITEA_HOST}. {@code HerdrPeerLauncher.applyGitToken} injects {@code GITEA_HOST}
* only for profiles that opted in via {@code gitTokenEnv} (CB-302), reading the value from
* that profile's {@code gitHostEnv} (default {@code GITEA_HOST}), so the shape only matters
* where a git token is opted in. A var used by more than one profile is one entry naming
* every profile that reads it, the same shape as {@link #requiredSecretEnvVars}.
*
* <p>Package-private and pure (no I/O, no logging) so the derivation is unit-testable
* without capturing log output; {@link #reportGitHostShape(FleetConfig)} is the logging caller.
*/
static Map<String, List<String>> gitHostEnvVars(FleetConfig cfg) {
Map<String, List<String>> hostsBy = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (profile.hasGitToken()) {
hostsBy.computeIfAbsent(profile.gitHostEnv(), _ -> new ArrayList<>())
.add("profile '" + name + "' gitHostEnv");
}
});
return hostsBy;
}
/**
* True when the value already starts with a URI scheme ({@code https://...}, {@code
* http://...}). A bare host name and a host:port must both report {@code false} — the shape
* this ticket exists for is a value that <em>looks like</em> a host but is a full URL, and
* confusing those in the report would move the failure to the log instead of the network.
*/
static boolean startsWithScheme(String value) {
return value.matches("[A-Za-z][A-Za-z0-9+.-]*://.*");
}
/**
* fleetd #377: log, on the same startup path as {@link #reportRequiredSecrets}, the SHAPE of
* each git host value that members receive as {@code GITEA_HOST}: set or unset, its length,
* whether it starts with a scheme, whether it ends with a slash. Never the value itself — the
* same discipline as {@link #reportRequiredSecrets}, which logs by name only. Either shape is
* legitimate: the value is passed to members unchanged, and a line that quietly rewrites it
* would change what works on one host and breaks on another. The shape line only tells the
* operator which URL form to expect from a member that builds on {@code GITEA_HOST}. An
* unset variable is logged at INFO — useful information, not an error — and the daemon
* starts on either way.
*
* <p>The env read lives in the overload below so a test can drive the line with a known
* value and prove that value never reaches the log.
*/
static void reportGitHostShape(FleetConfig cfg) {
reportGitHostShape(cfg, System.getenv());
}
static void reportGitHostShape(FleetConfig cfg, Map<String, String> env) {
Map<String, List<String>> hostsBy = gitHostEnvVars(cfg);
if (hostsBy.isEmpty()) {
log.info("startup git host: no profile sets a gitTokenEnv — nothing to check");
return;
}
hostsBy.forEach((varName, sources) -> {
String value = env.get(varName);
if (value == null || value.isBlank()) {
log.info("startup git host {}: unset ({}) — a member gets GITEA_TOKEN but no "
+ "GITEA_HOST value", varName, String.join(", ", sources));
} else {
log.info("startup git host {}: set ({}) — length={}, startsWithScheme={}, "
+ "trailingSlash={}",
varName, String.join(", ", sources),
value.length(), startsWithScheme(value), value.endsWith("/"));
}
});
}
/**
* fleetd #184: state the member trust model at startup. Environment controls and worktrees do
* not make a sandbox when fleetd and its members use the same OS user. A separate herdr may
@@ -38,15 +38,46 @@ public final class FleetHealthMonitor {
*/
static final long ASK_LAPSE_RECHECK_DELAY_SECONDS = 120;
/**
* fleetd #386: {@code System.nanoTime()} (or whatever {@link #clock} is) does not advance while
* macOS sleeps, so a raw {@code nowNanos - lastActivityAtNanos} comparison freezes with the
* host and can never cross {@link #workingSuspectAfterNanos}. This is a second, wall-clock
* source used ONLY inside the stall check ({@link #stallElapsedNanos}) to detect and correct
* for that freeze. Nothing else in this class reads it — every other decision (readiness grace,
* the fault classification itself) stays exactly on {@link #clock}, as the ticket requires.
*/
private static final LongSupplier DEFAULT_REALTIME_CLOCK =
() -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis());
private final AgentControl agents;
private final Supplier<List<MemberSession>> roster;
private final MessageService messages;
private final ScheduledExecutorService scheduler;
private final LongSupplier clock;
private final LongSupplier realtimeClock;
private final long intervalSeconds;
private final long tickIntervalNanos;
private final long workingSuspectAfterNanos;
private final BiConsumer<String, String> failTarget;
private final Map<String, HealthPrior> priors = new HashMap<>();
/**
* fleetd #386 clock-drift bookkeeping. {@code haveClockBaseline}/{@code lastTickMonoNanos}/
* {@code lastTickRealNanos} track the previous tick's pair of readings so each new tick can
* measure how far the two clocks moved apart since then. {@code accumulatedDriftNanos} is the
* running total of every such divergence observed since this monitor started (never decreases —
* the monotonic clock can only lag real time, never lead it). {@code busyDriftBaselineNanos}/
* {@code busyBaselineActivityNanos} record, per target, the value of {@code accumulatedDriftNanos}
* at the moment this monitor first saw that target's CURRENT {@code lastActivityAtNanos} while
* BUSY — so {@link #stallElapsedNanos} adds back only the drift observed DURING this BUSY span,
* never drift from a sleep that happened before the member went busy. All five fields are touched
* only from {@code tick()}, like {@link #priors}.
*/
private boolean haveClockBaseline = false;
private long lastTickMonoNanos;
private long lastTickRealNanos;
private long accumulatedDriftNanos = 0;
private final Map<String, Long> busyDriftBaselineNanos = new HashMap<>();
private final Map<String, Long> busyBaselineActivityNanos = new HashMap<>();
/**
* The live classification per member, and the only one of this class's three maps that more
* than one scheduler task touches. {@code tick} writes it (and prunes it to the roster);
@@ -88,12 +119,29 @@ public final class FleetHealthMonitor {
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
long workingSuspectAfterSeconds, BiConsumer<String, String> failTarget) {
this(agents, roster, messages, scheduler, clock, DEFAULT_REALTIME_CLOCK, intervalSeconds,
workingSuspectAfterSeconds, failTarget);
}
/**
* @param realtimeClock fleetd #386: a wall-clock nanosecond source (e.g.
* {@code System.currentTimeMillis()} converted to nanos) that keeps
* advancing while {@code clock} is frozen by a host sleep. Used only to
* correct the stall check — see the class-level javadoc on the
* clock-drift fields.
*/
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, LongSupplier realtimeClock,
long intervalSeconds, long workingSuspectAfterSeconds,
BiConsumer<String, String> failTarget) {
this.agents = agents;
this.roster = roster;
this.messages = messages;
this.scheduler = scheduler;
this.clock = clock;
this.realtimeClock = Objects.requireNonNull(realtimeClock, "realtimeClock");
this.intervalSeconds = intervalSeconds;
this.tickIntervalNanos = TimeUnit.SECONDS.toNanos(intervalSeconds);
this.workingSuspectAfterNanos = TimeUnit.SECONDS.toNanos(workingSuspectAfterSeconds);
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
}
@@ -123,6 +171,7 @@ public final class FleetHealthMonitor {
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
HashSet<String> current = new HashSet<>();
long nowNanos = clock.getAsLong();
long driftBeforeThisTick = observeClockDrift(nowNanos);
for (MemberSession session : rosterNow) {
current.add(session.terminalId());
Agent agent = live.get(session.terminalId());
@@ -133,7 +182,7 @@ public final class FleetHealthMonitor {
&& session.state() != MemberSession.State.SPAWNING;
boolean readinessGraceElapsed = nowNanos - session.spawnedAtNanos() >= READINESS_GRACE_NANOS;
boolean stalled = session.state() == MemberSession.State.BUSY
&& nowNanos - session.lastActivityAtNanos() >= workingSuspectAfterNanos;
&& stallElapsedNanos(session, nowNanos, driftBeforeThisTick) >= workingSuspectAfterNanos;
// CB-643: the three message-layer facts CB-640 published. Read them here rather than
// leaving them false — that constant is what made 8 of the 9 fault states dead.
boolean queuedDelivery = messages.hasQueuedDelivery(session.terminalId());
@@ -151,6 +200,8 @@ public final class FleetHealthMonitor {
priors.keySet().retainAll(current);
states.keySet().retainAll(current);
orphanStreaks.keySet().retainAll(current);
busyDriftBaselineNanos.keySet().retainAll(current);
busyBaselineActivityNanos.keySet().retainAll(current);
} catch (Throwable error) {
// Any unclassified collection failure must never kill the monitor's only scheduler task.
log.warn("fleet health collection failed; will retry next tick", error);
@@ -175,6 +226,65 @@ public final class FleetHealthMonitor {
return streak >= ORPHAN_CONFIRM_TICKS;
}
/**
* fleetd #386: compare this tick's monotonic and real-time readings against the previous
* tick's, and fold any positive divergence into {@link #accumulatedDriftNanos} (a ratchet — it
* never decreases, since the monotonic clock can only fall behind real time, never ahead of
* it). Logs once, at WARN, when that single tick's divergence exceeds one full tick interval —
* the signature of a host that slept between the two ticks (a tick literally cannot run while
* the process itself is suspended, so the whole sleep duration lands inside one tick's gap).
*
* @return {@link #accumulatedDriftNanos} as it stood BEFORE this tick's divergence was folded
* in — the baseline {@link #stallElapsedNanos} needs when a target is observed BUSY
* for the first time this tick, so a sleep that happened before this member went busy
* is not attributed to it.
*/
private long observeClockDrift(long nowNanos) {
long nowRealNanos = realtimeClock.getAsLong();
long driftBeforeThisTick = accumulatedDriftNanos;
if (haveClockBaseline) {
long monoDelta = nowNanos - lastTickMonoNanos;
long realDelta = nowRealNanos - lastTickRealNanos;
long tickDrift = realDelta - monoDelta;
if (tickDrift > tickIntervalNanos) {
log.warn("fleet health: the monotonic clock did not advance for about {}s that the "
+ "real clock did since the last tick (host likely slept); the stall "
+ "detector could not see that time", TimeUnit.NANOSECONDS.toSeconds(tickDrift));
}
if (tickDrift > 0) {
accumulatedDriftNanos = driftBeforeThisTick + tickDrift;
}
}
lastTickMonoNanos = nowNanos;
lastTickRealNanos = nowRealNanos;
haveClockBaseline = true;
return driftBeforeThisTick;
}
/**
* fleetd #386: {@code nowNanos - lastActivityAtNanos} alone freezes across a host sleep, since
* both come from the monotonic {@link #clock}. This adds back the real-time drift observed
* since this BUSY span started — not the monitor's whole lifetime, so a sleep that happened
* before this member went busy never leaks into its stall reading (see the class-level javadoc
* on the drift fields). The baseline resets whenever {@code lastActivityAtNanos} changes (a new
* turn) or the member is not currently BUSY.
*/
private long stallElapsedNanos(MemberSession session, long nowNanos, long driftBeforeThisTick) {
String target = session.terminalId();
if (session.state() != MemberSession.State.BUSY) {
busyDriftBaselineNanos.remove(target);
busyBaselineActivityNanos.remove(target);
return nowNanos - session.lastActivityAtNanos();
}
Long baselineActivity = busyBaselineActivityNanos.get(target);
if (baselineActivity == null || baselineActivity != session.lastActivityAtNanos()) {
busyBaselineActivityNanos.put(target, session.lastActivityAtNanos());
busyDriftBaselineNanos.put(target, driftBeforeThisTick);
}
long driftSinceBusyStart = accumulatedDriftNanos - busyDriftBaselineNanos.get(target);
return (nowNanos - session.lastActivityAtNanos()) + driftSinceBusyStart;
}
void reportTransition(String target, HealthState next) {
HealthState previous = states.put(target, next);
if (previous == next) return;
@@ -548,8 +548,19 @@ public final class CompletionResolver implements TurnListener {
* fast completion. Runs the same backend-error classification the normal and raw-scrape paths
* apply, against whatever is on screen right now: a match together with the too-fast crash
* signature notifies {@link #backendErrorSink} (only on the resolution that wins the race). A
* non-match stays the original generic too-fast failure, naming the member and both timings,
* with whatever the pane shows appended so the caller sees the cause, not just "it failed".
* non-match stays the generic too-fast failure, naming the member and both timings, with
* whatever the pane shows appended so the caller sees the evidence, not just "it failed".
*
* <p>fleetd#376: <strong>this path must never resolve a completion.</strong> A fast backend can
* genuinely answer inside the floor, so the failure is sometimes wrong — but it is wrong in the
* loud direction, and the fix for that is honest wording, not a guess at the pane's meaning.
* Reclassifying from the scrape was tried and rejected: there is no reliable positive marker for
* "this is a real reply" across backends. {@link #lastAssistantBlock} falls back to the entire
* pane when it finds no {@code ⏺} marker, so on a crash the candidate "reply" is the whole
* screen; and {@code ⏺} itself is a Claude Code marker that an opencode pane never carries — the
* very backend whose speed raised this ticket. Any weaker test (non-blank, or "contains sentence
* punctuation") passes on almost every crash, because a pane holding a file path, a version
* number or a hostname contains a full stop. That trades a loud wrong answer for a silent one.
*/
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
long elapsedNanos) {
@@ -560,13 +571,13 @@ public final class CompletionResolver implements TurnListener {
scrape = "";
}
String clippedScrape = clip(scrape);
String baseReason = String.format(
"member %s went BUSY -> DONE in %dms (floor %dms) — too fast to be real work, most "
+ "likely a backend error before any work started",
String timing = String.format(
"member %s went BUSY -> DONE in %dms (floor %dms)",
target, elapsedNanos / 1_000_000, MIN_TURN_NANOS / 1_000_000);
String backendError = firstMatchingLine(scrape, backendErrorPatternOrFallback(target));
if (backendError != null) {
String reason = baseReason + ": " + clippedScrape;
String reason = timing + " — too fast to be real work, and the pane carries a backend "
+ "error: " + clippedScrape;
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
@@ -575,7 +586,14 @@ public final class CompletionResolver implements TurnListener {
}
return;
}
fail(target, turn, clippedScrape.isBlank() ? baseReason : baseReason + ": " + clippedScrape);
// fleetd#376: no pattern matched, so the cause is genuinely unknown. Say that, rather than
// asserting a backend error the way this message used to — a fast backend really can finish
// inside the floor, and a reader who trusts a wrong cause stops looking at the pane.
String reason = timing + " — inside the floor. That is usually a backend error before any "
+ "work started, but a fast backend can answer inside it too, and nothing here can "
+ "tell those apart, so the turn is reported failed rather than guessed. Read the "
+ "pane below before deciding which it was";
fail(target, turn, clippedScrape.isBlank() ? reason : reason + ": " + clippedScrape);
}
/**
@@ -870,6 +870,11 @@ public final class FleetMcp {
* or — when no send is open — queueing the reply in the inbox for later drain (CB-307).
* {@code callerTerminal} is resolved from the connection (never an argument); a {@code null}
* means the caller is not a known worker (e.g. the primary called it by mistake).
*
* <p>fleetd #365: the result text names which of those actually happened
* ({@link MessageService.ReplyOutcome#description()}) instead of the single word "delivered"
* for both — a queued reply is a real success, but it is not the same fact as one that resolved
* a live waiter, and the caller could not previously tell them apart.
*/
static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) {
if (callerTerminal == null) {
@@ -884,8 +889,8 @@ public final class FleetMcp {
if (isBlank(content)) {
return error("content is required");
}
messages.reply(callerTerminal, content);
return text("delivered");
MessageService.ReplyOutcome outcome = messages.reply(callerTerminal, content);
return text(outcome.description());
}
/** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */
@@ -369,8 +369,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
// CB-547a: route the chosen profile but keep the caller's session identity — dropping it
// here would silently sever the resume handle on every policy-routed spawn. CB-557: the
// role rides along for the same reason, or a routed spawn would be labelled as a dev.
SpawnRequest routedReq = new SpawnRequest(chosen.profile(), req.requestedCwd(), req.callerCwd(),
req.sessionName(), req.resumeSessionId(), req.role());
SpawnRequest routedReq = req.withProfile(chosen.profile());
try {
PeerHandle handle = d.spawn(routedReq);
spawnedBy.put(handle.id(), d);
@@ -13,6 +13,7 @@ import java.nio.file.attribute.PosixFilePermissions;
import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.stream.Stream;
@@ -28,24 +29,46 @@ import java.util.stream.Stream;
* control lacked: herdr applies that overlay BEFORE the shell starts, so any sourced file can undo
* it — and did.
*
* <p><b>Which file is last depends on the platform, so the scrub runs from two of them.</b> zsh
* <p><b>zsh reads its four startup files under three different conditions, so no single file is
* guaranteed to run — the scrub has to cover the gap between them, not just the platforms.</b> zsh
* reads {@code .zshenv} always, {@code .zprofile} and {@code .zlogin} only for a LOGIN shell, and
* {@code .zshrc} only for an INTERACTIVE one. herdr does not open the same kind of shell
* everywhere — measured on herdr 0.8.0: macOS panes run {@code -zsh} (login, so {@code .zlogin}
* runs), Linux panes run a plain {@code /usr/bin/zsh} (interactive but NOT login, so
* {@code .zlogin} never runs at all). A scrub in {@code .zlogin} alone is therefore a control that
* silently does nothing on Linux — the exact failure this class exists to remove, one platform
* over.
* {@code .zshrc} only for an INTERACTIVE one. A pane shell that is at least one of login or
* interactive is covered by sourcing the scrub from {@code .zshrc} and {@code .zlogin} (below), but
* a pane shell that is NEITHER reads only {@code .zshenv} and stops — fleetd #388, measured: a herdr
* pane can be neither login nor interactive, and such a pane read {@code .zshenv}, never reached
* {@code scrub.zsh}, and left no report at all. A bare {@code argv[0]} of {@code /usr/bin/zsh}
* proves the shell is NOT a login shell; it says nothing about whether it is interactive, so it
* must never be read as "therefore interactive" — that wrong inference is what let #388 ship.
*
* <p>So both {@code .zshrc} and {@code .zlogin} source the same generated {@code scrub.zsh} after
* sourcing their {@code $HOME} counterpart. On Linux only the first fires; on macOS both do, and
* the second pass is deliberate rather than merely harmless — it re-scrubs anything the operator's
* own {@code ~/.zlogin} exported after {@code .zshrc} had finished. Re-running is idempotent: a
* name already blank is blanked again, and the report is rewritten with the same counts.
* <p>So {@code .zshenv} carries a THIRD pass, guarded by the exact condition that defines the gap:
* {@code [[ ! -o login && ! -o interactive ]]}. That guard is why this pass cannot double-scrub a
* pane that {@code .zshrc} or {@code .zlogin} will also cover — one of {@code -o login}/
* {@code -o interactive} is always true there, so the {@code .zshenv} pass never fires for them, and
* their own unconditional sourcing is untouched. The guard also carries a sentinel
* ({@value #SCRUB_SENTINEL}) so it fires once per PANE and not once per PROCESS: {@code .zshenv} is
* read by every zsh a member's own tooling forks (a plain {@code zsh -c '...'} for a single
* command is itself neither login nor interactive), and those children inherit variables their
* parent deliberately set for them (git hooks get {@code GIT_DIR}, a venv gets
* {@code VIRTUAL_ENV}, a build tool gets {@code NODE_OPTIONS} or {@code JAVA_TOOL_OPTIONS}).
* Re-scrubbing every such child would blank all of that, and would also make the pane's own
* {@code scrub-report.txt} — rewritten on every pass — describe whichever child exited last
* instead of the pane. The sentinel is exported only AFTER {@code scrub.zsh} runs, so the pass
* that sets it never sees it and cannot blank it; it must also be on the scrub's own allow-list
* (see {@link #generate(Path, Set)}) so a later pass, in the same pane, cannot blank it back to
* empty — an exported-but-empty sentinel reads as unset to the {@code -z} guard and would silently
* re-enable scrubbing for every subsequent child of that pane.
*
* <p>So all three of {@code .zshenv} (gap only, guarded), {@code .zshrc}, and {@code .zlogin}
* source the same generated {@code scrub.zsh} after sourcing their {@code $HOME} counterpart. A
* login-and-interactive pane runs the {@code .zshrc} and {@code .zlogin} passes, and the second is
* deliberate rather than merely harmless — it re-scrubs anything the operator's own
* {@code ~/.zlogin} exported after {@code .zshrc} had finished. A pane that is neither runs only the
* {@code .zshenv} pass. Re-running is idempotent: a name already blank is blanked again, and the
* report is rewritten with the same counts.
*
* <p>Each generated file sources its {@code $HOME} counterpart FIRST, so {@code PATH} and every
* toolchain binary still resolve exactly as the operator configured them; only afterwards does
* {@code .zlogin} run the scrub: every EXPORTED variable not on the derived allow-list is re-exported
* toolchain binary still resolve exactly as the operator configured them; only afterwards does the
* scrub run: every EXPORTED variable not on the derived allow-list is re-exported
* blank. Blank, not credential-shaped-pattern-filtered: a pattern list ({@code *TOKEN*}, …) is an
* enumeration and misses what it did not think of — a username is the other half of a credential and
* is shaped like none. Credential-SHAPED names among the blanked set go to the WARN log only,
@@ -63,7 +86,11 @@ public final class EnvAllowListScrub {
/** Name of the report file written into the generated directory by the scrub itself. */
static final String REPORT_FILE = "scrub-report.txt";
/** The scrub body, generated once and sourced from both {@code .zshrc} and {@code .zlogin}. */
/**
* The scrub body, generated once and sourced from {@code .zshrc} and {@code .zlogin}
* unconditionally, and from {@code .zshenv} when the pane shell is neither login nor
* interactive (fleetd #388) — see the class javadoc.
*/
static final String SCRUB_FILE = "scrub.zsh";
/** Prefix of every generated directory — also what {@link #reapOrphans} matches on. */
@@ -79,6 +106,28 @@ public final class EnvAllowListScrub {
private static final String SOURCE_SCRUB =
"source \"$ZDOTDIR/" + SCRUB_FILE + "\"\n";
/**
* fleetd #388: marks a pane, not a process, as already scrubbed. Set only by the guarded
* {@code .zshenv} pass (see {@link #NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB}) after
* {@code scrub.zsh} has run, so it must also be folded into that pass's own allow-list — see
* the class javadoc's "must also be on the scrub's own allow-list" paragraph.
*/
static final String SCRUB_SENTINEL = "_CB633_SCRUBBED";
/**
* Appended to {@code .zshenv}, after its {@code $HOME} source: the third pass, guarded on the
* exact condition that defines the gap {@code .zshrc}/{@code .zlogin} do not cover — a shell
* that is neither login nor interactive. The sentinel export happens only once the scrub has
* already run, and only for as long as the current pane's environment has not been rebuilt from
* scratch (a fresh {@code env -i} child would not inherit it — that is out of scope here, since
* such a child is no longer running under the pane's own environment at all).
*/
private static final String NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB =
"if [[ ! -o login && ! -o interactive && -z \"${" + SCRUB_SENTINEL + ":-}\" ]]; then\n"
+ " " + SOURCE_SCRUB
+ " export " + SCRUB_SENTINEL + "=1\n"
+ "fi\n";
private EnvAllowListScrub() {
}
@@ -105,11 +154,16 @@ public final class EnvAllowListScrub {
reapOrphans(parentDir);
Path dir = Files.createTempDirectory(parentDir, DIR_PREFIX);
dir.toFile().deleteOnExit();
// The report is written by zsh, after these hooks are registered, so register its path
// too — otherwise the directory is non-empty at JVM exit and cannot be removed at all.
dir.resolve(REPORT_FILE).toFile().deleteOnExit();
write(dir, SCRUB_FILE, scrubScript(allowedNames));
write(dir, ".zshenv", homeSourcingFile(".zshenv"));
// zsh truncates this pre-created receipt after these hooks are registered. Register its
// path too — otherwise the directory is non-empty at JVM exit and cannot be removed.
Files.createFile(dir.resolve(REPORT_FILE)).toFile().deleteOnExit();
// fleetd #388: scrub.zsh's OWN allow-list must also keep SCRUB_SENTINEL, or a later
// pass in the same pane blanks it back to empty and the .zshenv guard below thinks it
// was never scrubbed — see the class javadoc.
Set<String> namesForScrubScript = new HashSet<>(allowedNames);
namesForScrubScript.add(SCRUB_SENTINEL);
write(dir, SCRUB_FILE, scrubScript(namesForScrubScript));
write(dir, ".zshenv", homeSourcingFile(".zshenv") + NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB);
write(dir, ".zprofile", homeSourcingFile(".zprofile"));
write(dir, ".zshrc", homeSourcingFile(".zshrc") + SOURCE_SCRUB);
write(dir, ".zlogin", homeSourcingFile(".zlogin") + SOURCE_SCRUB);
@@ -147,12 +201,11 @@ public final class EnvAllowListScrub {
* into it: owner keeps full access, {@code group} gets traverse+read on the directory ({@code
* rwxr-x---}, so a member process — a login shell reading it via {@code ZDOTDIR}, or another
* process simply opening a file under it — running under that group can find and read the
* files) and read-only on each file ({@code rw-r-----}) — deliberately no group WRITE anywhere,
* since a member never needs to add or change what fleetd generated. (For the ZDOTDIR scrub
* specifically, this also means the scrub script's own report write inside the pane fails
* closed rather than open — see {@code scrub.zsh}'s trailing {@code 2>/dev/null} — which {@link
* dev.ltms.fleet.member.HerdrPeerLauncher#releaseZdotdir} already treats as "cannot be
* confirmed to have run" rather than success.)
* files) and read-only on each file ({@code rw-r-----}), except the pre-created ZDOTDIR
* {@code scrub-report.txt}. That receipt gets group write ({@code rw-rw----}), so
* {@code scrub.zsh} can truncate and write it without granting group write on the directory.
* If its optional permission change fails, the member cannot write a receipt and the launcher
* keeps its existing WARN rather than failing the spawn.
*
* <p>Package-private and named generically on purpose: fleetd #213 built this for the ZDOTDIR
* scrub directory, and fleetd #219 reuses it verbatim for {@link
@@ -167,7 +220,15 @@ public final class EnvAllowListScrub {
setGroupAndPermissions(dir, principal, "rwxr-x---");
try (Stream<Path> entries = Files.list(dir)) {
for (Path file : entries.toList()) {
setGroupAndPermissions(file, principal, "rw-r-----");
if (REPORT_FILE.equals(file.getFileName().toString())) {
try {
setGroupAndPermissions(file, principal, "rw-rw----");
} catch (IOException | UnsupportedOperationException ignored) {
// The receipt is optional. Its absence keeps the existing WARN path.
}
} else {
setGroupAndPermissions(file, principal, "rw-r-----");
}
}
}
} catch (IOException e) {
@@ -245,11 +306,13 @@ public final class EnvAllowListScrub {
}
return """
# generated by fleetd (CB-633 memberCredentials policy=allow-list) — do not edit.
# Sourced from .zshrc and again from .zlogin, each time AFTER that file has sourced
# its $HOME counterpart — so this runs after everything the operator sourced, on a
# login shell (macOS panes) and on a plain interactive one (Linux panes) alike.
# Running twice is idempotent and deliberate: the second pass catches anything
# ~/.zlogin exported after ~/.zshrc had finished.
# Sourced unconditionally from .zshrc and again from .zlogin, each time AFTER that
# file has sourced its $HOME counterpart — so this runs after everything the
# operator sourced, on any pane that is login and/or interactive. Running twice is
# idempotent and deliberate: the second pass catches anything ~/.zlogin exported
# after ~/.zshrc had finished. Also sourced, once, from a guarded pass in .zshenv
# (fleetd #388) when the pane shell is NEITHER login nor interactive — the one gap
# those two files do not cover.
typeset -A _cb633_allowed
for _cb633_n in %s; do _cb633_allowed[$_cb633_n]=1; done
@@ -56,10 +56,13 @@ public final class FleetMetrics {
m.describe(REPLIES, "counter",
"Worker replies by delivery path (rendezvous=resolved an open send, inbox=stranded and held).");
m.describe(PUSH_NUDGES, "counter",
"CB-307 push-loop nudges to the primary (delivered|exhausted).");
"CB-307 push-loop nudges to the primary (sent|exhausted). fleetd #365: \"sent\" means "
+ "the herdr paste-and-submit call succeeded, not that the pane read it — this "
+ "layer has no read-receipt concept.");
m.describe(HEARTBEAT_NUDGES, "counter",
"CB-551 idle-lead heartbeat nudges (delivered|failed|exhausted). Quiet-cap exhaustion "
+ "means the lead idled with nothing pending and was told to stand down.");
"CB-551 idle-lead heartbeat nudges (sent|failed|exhausted). Quiet-cap exhaustion "
+ "means the lead idled with nothing pending and was told to stand down. "
+ "fleetd #365: \"sent\" means the herdr call succeeded, not that the lead read it.");
m.describe(SPAWNS, "counter",
"Worker spawn attempts by peer kind and outcome (ready|timeout|guard_rejected).");
m.describe(HERDR_CALLS, "counter",
@@ -32,7 +32,11 @@ public interface LeadChannel {
/** Non-destructive FIFO snapshot of the messages held for this daemon's own coord-id. */
List<LeadMessage> peek();
/** Drop {@code msgId} from the held set and ack it on the broker. A no-op if it is not held. */
/**
* Drop {@code msgId} from the held set and ack it on the broker. A repeated ack that this
* connection already completed may be a no-op. Any other unknown msgId must throw rather than
* report an ack that did not reach the broker.
*/
void ack(String msgId);
/** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
@@ -5,6 +5,7 @@ import dev.ltms.fleet.herdr.AgentStatus;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ScheduledExecutorService;
@@ -43,6 +44,9 @@ public final class LeadCoordLoop {
private static final Logger log = LoggerFactory.getLogger(LeadCoordLoop.class);
/** A bounded window is enough: redelivery can only follow a recent failed ack or recovery. */
private static final int RECENT_DELIVERY_LIMIT = 1_024;
/** How an arriving peer message is rendered into the lead's pane — the sender's coord-id, then its text. */
static final String DELIVERY_FORMAT = "[lead %s] %s";
@@ -51,6 +55,8 @@ public final class LeadCoordLoop {
private final Supplier<Map<String, String>> leads;
private final ScheduledExecutorService scheduler;
private final long intervalMs;
/** msgIds already written to the pane, so recovery redelivery is acked without another pane write. */
private final LinkedHashMap<String, Boolean> delivered = new LinkedHashMap<>();
private volatile boolean running;
@@ -117,6 +123,13 @@ public final class LeadCoordLoop {
if (held.isEmpty()) {
return;
}
LeadMessage msg = held.getFirst();
if (wasDelivered(msg.msgId())) {
// This lives here, rather than in LeadMailbox, because only this loop knows a pane write
// happened. The mailbox only knows broker delivery tags and must still redeliver after a crash.
ackDelivered(msg);
return;
}
String lead = resolveLocalLead();
if (lead == null) {
// Left unacked on purpose: the broker keeps holding it until a lead pane exists.
@@ -136,7 +149,6 @@ public final class LeadCoordLoop {
lead, status, held.size());
return;
}
LeadMessage msg = held.getFirst();
try {
agents.send(lead, DELIVERY_FORMAT.formatted(msg.from(), msg.content()));
} catch (RuntimeException e) {
@@ -145,6 +157,13 @@ public final class LeadCoordLoop {
msg.msgId(), msg.from(), lead, e.toString());
return;
}
rememberDelivered(msg.msgId());
if (ackDelivered(msg)) {
log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead);
}
}
private boolean ackDelivered(LeadMessage msg) {
try {
channel.ack(msg.msgId());
} catch (RuntimeException e) {
@@ -152,9 +171,24 @@ public final class LeadCoordLoop {
// deliberate direction of this trade.
log.warn("lead coordination: delivered message {} but could not ack it: {}",
msg.msgId(), e.toString());
return;
return false;
}
return true;
}
private boolean wasDelivered(String msgId) {
synchronized (delivered) {
return delivered.containsKey(msgId);
}
}
private void rememberDelivered(String msgId) {
synchronized (delivered) {
delivered.put(msgId, Boolean.TRUE);
if (delivered.size() > RECENT_DELIVERY_LIMIT) {
delivered.remove(delivered.keySet().iterator().next());
}
}
log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead);
}
/**
@@ -96,7 +96,14 @@ public final class LeadHeartbeatLoop {
this.metrics = metrics;
}
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
/**
* Count one nudge outcome when a registry is wired; a no-op in unit tests.
*
* <p>fleetd #365: the {@code "sent"} outcome (renamed from {@code "delivered"}) records only
* that {@link #injectNudge} — a one-way herdr {@code agent.prompt} paste-and-submit — returned
* without throwing, not that the lead's pane actually read or acted on the text. This layer has
* no read-receipt concept, so "sent" is the honest word for what this call can ever establish.
*/
private void countNudge(String outcome) {
if (metrics != null) {
metrics.inc(FleetMetrics.HEARTBEAT_NUDGES, "outcome", outcome);
@@ -245,7 +252,7 @@ public final class LeadHeartbeatLoop {
agents.send(leadTerminal, fleet.nudgeText());
log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})",
leadTerminal, quietCount);
countNudge("delivered");
countNudge("sent");
} catch (RuntimeException e) {
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
countNudge("failed");
@@ -59,7 +59,8 @@ import java.util.concurrent.TimeoutException;
* <p><strong>Recovery.</strong> The connection is opened with automatic + topology recovery
* enabled, mirroring {@code AmqpReplyInbox}: on reconnect the broker hands out fresh delivery tags,
* so the held snapshot is cleared (dedup by {@code msgId} still prevents any double-queue on
* redelivery) and any publish still awaiting its confirm is failed rather than left to idle out
* redelivery). {@link LeadCoordLoop} separately deduplicates pane writes, since it alone knows
* which messages reached a lead. Any publish still awaiting its confirm is failed rather than left to idle out
* the confirm timeout against a sequence number that means nothing on the new channel.
*/
public final class LeadMailbox implements LeadChannel, AutoCloseable {
@@ -85,6 +86,10 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
private final Object channelLock = new Object();
/** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */
private final LinkedHashMap<String, Held> held = new LinkedHashMap<>();
/** Successful broker acks on this connection, retained only to make a repeated caller ack quiet. */
private final LinkedHashMap<String, Boolean> recentlyAcked = new LinkedHashMap<>();
/** Bounds {@link #recentlyAcked}: it is only an idempotency aid, never delivery state. */
private static final int RECENT_ACK_LIMIT = 1_024;
/**
* A dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + ack)
@@ -162,16 +167,15 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
}
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). LeadCoordLoop
// remembers successful pane writes separately, so that redelivery cannot write a pane twice. Any publish
// confirm still in flight when the connection dropped is equally stale — fail it now rather
// than let it silently ride out CONFIRM_TIMEOUT_MS.
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
synchronized (held) {
held.clear();
}
clearHeldForRecovery();
failPendingPublishesOnRecovery();
log.info("AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery");
}
@@ -350,15 +354,26 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
return snapshot;
}
/** Remove the held message {@code msgId} and ack it on the broker. No-op if not held. */
/**
* Remove the held message {@code msgId} and ack it on the broker.
*
* <p>A repeated ack that this connection already completed is a no-op, tracked in the bounded
* {@link #recentlyAcked} set. Any other missing entry throws: recovery clears {@link #held} while
* the broker still owns the unacked delivery, and quiet success there would hide a required retry.
* The set is bounded because it only distinguishes a recent duplicate caller ack from an unknown
* delivery; it is not a substitute for broker state across a reconnect.
*/
@Override
public void ack(String msgId) {
Held h;
synchronized (held) {
h = held.remove(msgId);
if (h == null && recentlyAcked.containsKey(msgId)) {
return;
}
}
if (h == null) {
return; // never held (or already acked) — no-op
throw new IllegalStateException("cannot ack lead message " + msgId + ": it is not held");
}
try {
synchronized (channelLock) {
@@ -372,6 +387,19 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
}
throw new IllegalStateException("cannot ack lead message " + msgId, e);
}
synchronized (held) {
recentlyAcked.put(msgId, Boolean.TRUE);
if (recentlyAcked.size() > RECENT_ACK_LIMIT) {
recentlyAcked.remove(recentlyAcked.keySet().iterator().next());
}
}
}
/** Clear stale delivery tags after recovery; package-private so the recovery contract test drives this exact path. */
void clearHeldForRecovery() {
synchronized (held) {
held.clear();
}
}
private DeliverCallback deliverCallback() {
@@ -138,6 +138,53 @@ public final class MessageService {
public record AskResult(AskOutcome outcome, String answer) {
}
/**
* How a worker's {@code fleet_reply} ({@link #reply(String, String)}) actually landed
* (fleetd #365) — the two doors that expose it, {@code fleet_reply} and {@code POST
* /sessions/{id}/reply}, both used to report the single word "delivered" whichever of these
* happened, so a caller could not tell an active handoff from a reply merely held for later
* drain. Both are successes; they are not the same fact.
*/
public enum ReplyOutcome {
/** Resolved a {@code fleet_send}/{@code fleet_ask} that was actively waiting on this reply. */
RESOLVED_SEND("resolved_send", true,
"delivered — resolved the fleet_send that was waiting for it"),
/**
* No live waiter was open, but the reply completed a parked async ticket directly
* ({@link #askAnsweredAsyncTasks}) — a {@code fleet_poll} caller sees it immediately.
*/
RESOLVED_ASYNC_TICKET("resolved_async_ticket", true,
"delivered — resolved a pending async ticket (visible to fleet_poll)"),
/** Nothing was waiting; the reply was queued in the inbox for a later drain (CB-307). */
QUEUED("queued", false,
"queued — no send or ticket was waiting; held in the inbox for a later drain");
private final String wireName;
private final boolean delivered;
private final String description;
ReplyOutcome(String wireName, boolean delivered, String description) {
this.wireName = wireName;
this.delivered = delivered;
this.description = description;
}
/** Stable machine-readable name for a JSON/metrics label (REST's {@code outcome} field). */
public String wireName() {
return wireName;
}
/** Whether something was actively waiting and received this reply right now. */
public boolean delivered() {
return delivered;
}
/** Shared human-readable text — the one place both {@code fleet_reply} and REST word this. */
public String description() {
return description;
}
}
/** Lifecycle phase of an async delegation ticket. */
public enum Phase {
/** Delegated and in flight — queued for the worker or being worked. */
@@ -439,16 +486,15 @@ public final class MessageService {
* @throws IllegalArgumentException if {@code content} is {@code null} or blank — the caller must
* report this as a client error (REST: 400 {@code bad_request}) rather than resolve
* anything
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
* was queued
* @return which of the three ways (fleetd #365) the reply actually landed — never {@code null}
*/
public boolean reply(String session, String content) {
public ReplyOutcome reply(String session, String content) {
if (content == null || content.isBlank()) {
throw new IllegalArgumentException("content is required");
}
if (rendezvous.resolve(session, content)) {
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
return ReplyOutcome.RESOLVED_SEND; // a live send took it — unchanged fast path
}
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
@@ -480,7 +526,7 @@ public final class MessageService {
asyncTasksByTurn.remove(turnId, orphan);
}
count(FleetMetrics.REPLIES, "path", "async-recovered");
return true; // the ticket itself took it — no inbox stranding at all
return ReplyOutcome.RESOLVED_ASYNC_TICKET; // the ticket itself took it — no inbox stranding
}
} else if (candidates.size() > 1) {
List<String> tickets = candidates.stream().map(t -> t.ticket).toList();
@@ -498,7 +544,7 @@ public final class MessageService {
if (pushLoop != null) {
pushLoop.onReplyQueued(session);
}
return true; // held, not lost
return ReplyOutcome.QUEUED; // held, not lost — but not delivered either
}
/**
@@ -132,7 +132,15 @@ public final class ReplyPushLoop {
this.metrics = metrics;
}
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
/**
* Count one nudge outcome when a registry is wired; a no-op in unit tests.
*
* <p>fleetd #365: the {@code "sent"} outcome (renamed from {@code "delivered"}) records only
* that {@code agents.send} — a one-way herdr {@code agent.prompt} paste-and-submit — returned
* without throwing. Nothing in this loop, or anywhere downstream of it, confirms the pane
* actually read or acted on the text; there is no read-receipt concept at this layer. "Sent"
* says exactly that; "delivered" claimed more than this call can ever establish.
*/
private void countNudge(String outcome) {
if (metrics != null) {
metrics.inc(FleetMetrics.PUSH_NUDGES, "outcome", outcome);
@@ -730,7 +738,7 @@ public final class ReplyPushLoop {
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders,
replyTargets.size(), tickets.size(), questions.size());
countNudge("delivered");
countNudge("sent");
for (PendingIncident incident : incidents) {
if (pendingIncidents.remove(incident.key(), incident)) {
deliveredIncidents.add(incident.key());
@@ -29,4 +29,9 @@ public record SpawnRequest(String profileName, String requestedCwd, String calle
String sessionName, String resumeSessionId) {
this(profileName, requestedCwd, callerCwd, sessionName, resumeSessionId, null);
}
/** Return a copy of this request with {@code profileName} replaced by {@code profile}. */
public SpawnRequest withProfile(String profile) {
return new SpawnRequest(profile, requestedCwd, callerCwd, sessionName, resumeSessionId, role);
}
}
@@ -696,6 +696,10 @@ public final class FleetApp {
/**
* The worker's structured reply ({@code fleet_reply}) — resolves the blocking send awaiting
* on this session, or queues the reply in the inbox when no send is open (CB-307).
*
* <p>fleetd #365: the response body's {@code delivered} field used to be unconditionally
* {@code true} for either case; it now reports whether a send/ticket was actually resolved,
* with {@code outcome} naming which (see {@link MessageService.ReplyOutcome}).
*/
private void replyMessage(Context ctx) {
String id = ctx.pathParam("id");
@@ -719,13 +723,19 @@ public final class FleetApp {
// a WRONG value instead of failing loudly. The check lives in MessageService.reply so both
// this door and FleetMcp.reply inherit the same rule; this catch only translates it into the
// {error, detail} envelope this file uses everywhere else.
MessageService.ReplyOutcome outcome;
try {
messages.reply(id, content);
outcome = messages.reply(id, content);
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", e.getMessage()));
return;
}
ctx.status(200).json(Map.of("sessionId", id, "delivered", true));
// fleetd #365: "delivered": true used to be unconditional here, whether the reply resolved
// a waiting send or was merely queued in the inbox for a later drain — the same gap
// FleetMcp.reply had over MCP. `delivered` now reflects which actually happened, and
// `outcome` names the specific case (see MessageService.ReplyOutcome).
ctx.status(200).json(Map.of("sessionId", id, "delivered", outcome.delivered(),
"outcome", outcome.wireName()));
}
/**
@@ -0,0 +1,241 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Level;
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;
import java.util.List;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #377: the startup git host line must report the <em>shape</em> of the value a
* member receives as GITEA_HOST — set or unset, length, scheme, trailing slash — and the
* value itself must never reach the log.
*/
class GitHostShapeReportTest {
/** A profile that opts in to the git-forge token, so GITEA_HOST is the var that matters. */
private static final String GIT_TOKEN_CONFIG = """
profiles:
local:
baseUrl: http://gx00.gw:8000
gitTokenEnv: WORKER_GITEA_TOKEN
""";
private static FleetConfig load(Path dir, String yaml) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml);
return FleetConfig.load(f);
}
/**
* The level this logger had before {@link #attach()} raised it, so {@link #detach} can put it
* back. {@code null} is a real value here — it means "inherit from the parent" — and that is
* exactly the state this logger starts in, so it must be restored as {@code null} rather than
* as some concrete level.
*/
private static Level originalLevel;
/**
* fleetd #377: the shape lines are logged at INFO, and {@code logback-test.xml} sets
* {@code dev.ltms.fleet} to WARN — so INFO events are dropped by the level check BEFORE any
* appender sees them. Attaching an appender is therefore not enough: without raising the level
* the list stays empty and every assertion below fails against correct production code. The
* sibling report tests do the same thing at each call site (see
* {@code MemberTrustModelReportTest}); doing it here keeps it in one place.
*/
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
originalLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
logger.detachAppender(appender);
logger.setLevel(originalLevel);
}
private static List<String> messages(ListAppender<ILoggingEvent> appender) {
return appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
}
@Test
void setLineReportsLengthSchemeAndTrailingSlash(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
String value = "https://git.example.test/";
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
} finally {
detach(appender);
}
String expected = "startup git host GITEA_HOST: set (profile 'local' gitHostEnv) — "
+ "length=" + value.length() + ", startsWithScheme=true, trailingSlash=true";
assertTrue(messages(appender).contains(expected),
"a value with a scheme and a trailing slash must be reported by shape only: "
+ "its length, startsWithScheme=true, trailingSlash=true");
}
@Test
void bareHostLineReportsNoSchemeNoTrailingSlash(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
String value = "git.example.test";
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
} finally {
detach(appender);
}
String expected = "startup git host GITEA_HOST: set (profile 'local' gitHostEnv) — "
+ "length=" + value.length() + ", startsWithScheme=false, trailingSlash=false";
assertTrue(messages(appender).contains(expected),
"a bare host name must report startsWithScheme=false, trailingSlash=false");
}
@Test
void unsetVariableIsLoggedAtInfoAndDoesNotThrow(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
ListAppender<ILoggingEvent> appender = attach();
try {
// GITEA_HOST absent from the map entirely — the unset case must not throw.
Fleetd.reportGitHostShape(cfg, Map.of());
} finally {
detach(appender);
}
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.INFO
&& e.getFormattedMessage()
.startsWith("startup git host GITEA_HOST: unset (profile 'local' gitHostEnv)")),
"an unset git host is useful information, not an error — say it at INFO level");
}
/**
* The important test: the value must NEVER appear in the log output. This test fails if the
* line is ever changed to include the value, because it drives the real logging path with a
* value that carries a marker no shape field could contain.
*/
@Test
void theValueNeverAppearsInLogOutput(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
String marker = "never-logged-host-shape-377";
String value = "https://" + marker + "/";
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
} finally {
detach(appender);
}
List<String> msgs = messages(appender);
assertFalse(msgs.stream().anyMatch(m -> m.contains(value)),
"the full GITEA_HOST value must never reach the log");
assertFalse(msgs.stream().anyMatch(m -> m.contains(marker)),
"no fragment of the value may reach the log — a line that embeds the value"
+ " would leak at least this marker");
assertTrue(msgs.stream().anyMatch(m -> m.contains("startsWithScheme=true")
&& m.contains("trailingSlash=true")),
"the shape must still be reported although the value is not");
}
@Test
void unsetReportAlsoCarriesNoValue(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
// A blank value is not injected either (putIfPresent skips it), so it must report unset
// without echoing anything of it.
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", " "));
} finally {
detach(appender);
}
assertTrue(messages(appender).stream().anyMatch(m ->
m.startsWith("startup git host GITEA_HOST: unset")),
"a blank value is skipped by the launcher, so the line reports unset");
}
@Test
void defaultGitHostEnvIsGITEA_HOST(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
assertEquals(Map.of("GITEA_HOST", List.of("profile 'local' gitHostEnv")),
Fleetd.gitHostEnvVars(cfg));
}
@Test
void explicitGitHostEnvReportsUnderItsOwnName(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
gitTokenEnv: WORKER_GITEA_TOKEN
gitHostEnv: MY_FORGE_HOST
""");
String value = "https://forge.example.test/";
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("MY_FORGE_HOST", value));
} finally {
detach(appender);
}
assertTrue(messages(appender).contains(
"startup git host MY_FORGE_HOST: set (profile 'local' gitHostEnv) — "
+ "length=" + value.length()
+ ", startsWithScheme=true, trailingSlash=true"),
"an explicitly named gitHostEnv is reported under that name");
}
@Test
void noGitTokenMeansNoHostLineIsNeeded(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of());
} finally {
detach(appender);
}
assertEquals(List.of("startup git host: no profile sets a gitTokenEnv — nothing to check"),
messages(appender));
}
@Test
void aHostWithAPortIsNotTreatedAsAScheme() {
assertFalse(Fleetd.startsWithScheme("git.example.test"));
assertFalse(Fleetd.startsWithScheme("git.example.test:3000"));
assertFalse(Fleetd.startsWithScheme("git.example.test:3000/"));
assertTrue(Fleetd.startsWithScheme("https://git.example.test"));
assertTrue(Fleetd.startsWithScheme("https://git.example.test:3000/"));
assertTrue(Fleetd.startsWithScheme("ssh://git.example.test"));
}
}
@@ -0,0 +1,172 @@
package dev.ltms.fleet.config;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Fleetd #358, the same "defect factory" #357 guarded on {@code FleetConfig.withDefaults()}
* (see {@code FleetConfigWithDefaultsPreservesEveryComponentTest}), reproduced here on
* {@link FleetConfig.Profile}. {@code Profile} carries a long back-compat constructor ladder — 8
* constructors, re-counted directly against the source rather than trusted from the ticket, at
* arities 25, 24, 22, 20, 18, 15, 14 and 12, against a canonical arity of 26 — and exactly ONE
* rebuild site, {@link FleetConfig.Profile#withProfile(String)}, whose own
* {@code return new Profile(...)} call is written at a literal 26-arg count. Add a 27th component
* and its established back-compat constructor at the old (26-arg) arity, and {@code withProfile}'s
* own call becomes a legal match for that new overload — silently dropping the new component every
* time a profile's name is defaulted from its {@code workers:} key.
*
* <p>Builds one {@link FleetConfig.Profile} through the TRUE canonical constructor — resolved by
* the record's own component types via {@code getDeclaredConstructor}, never by argument count —
* with a real, distinctive, non-null value in every component, calls {@link
* FleetConfig.Profile#withProfile(String)}, and asserts every component except {@code profile}
* itself survives unchanged, while {@code profile} comes back as the new name it was given.
*
* <p>Every value here is chosen so {@code Profile}'s own compact constructor (which normalizes
* several components — defaults {@code argv}/{@code kind}/{@code placement}/{@code workspace}/
* {@code gitHostEnv}, nulls a handful of blank-checked strings, clamps {@code weight}, coerces
* {@code subscription}) leaves it unchanged: every String is non-blank and already in the shape the
* compact constructor would otherwise coerce it to (e.g. {@code placement} is already lowercase),
* and every collection is non-empty. That is what makes "must survive unchanged" a valid assertion
* for every component below, the same reasoning {@code FleetConfigWithDefaultsPreservesEveryComponentTest}
* documents for {@code withDefaults()}.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
* {@link #exclusionListSizeIsPinned()} — a checker whose escape hatch can grow to silence a failure
* is not a checker. Every one of {@code Profile}'s 26 current components has a real, non-null,
* non-blank value here and none is excluded.
*/
class FleetConfigProfileWithProfilePreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.Profile.class.getRecordComponents();
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null value per component, chosen to survive the compact ctor. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profile", "profile-guard");
v.put("baseUrl", "https://guard.example/base");
v.put("model", "model-guard");
v.put("configDir", "/config/guard");
v.put("tokenEnv", "GUARD_TOKEN");
v.put("argv", List.of("guard-cmd"));
v.put("placement", "guard-placement");
v.put("workspace", "workspace-guard");
v.put("tabLabel", "tab-guard");
v.put("mcpUrl", "https://mcp.guard/");
v.put("cwd", "/cwd/guard");
v.put("parityOverlay", List.of(".guardrc"));
v.put("gitTokenEnv", "GUARD_GIT_TOKEN");
v.put("gitHostEnv", "GUARD_GIT_HOST");
v.put("kind", "claude-code");
v.put("env", Map.of("GUARD_ENV", "1"));
v.put("weight", 2.5f);
v.put("maxLoad", 4);
v.put("subscription", Boolean.TRUE);
v.put("exhaustedPattern", "pattern-guard");
v.put("credentialId", "cred-guard");
v.put("ideMcpUrl", "https://ide.guard/");
v.put("ideProjectDir", "ide-project-guard");
v.put("ideOpenCommand", "open-guard {dir}");
v.put("autoCompactWindow", 150_000);
v.put("errorPattern", "error-pattern-guard");
assertNamesMatchComponents(v);
return v;
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
* to add a new component here fails this assertion by name, rather than silently checking one
* component fewer than the record has.
*/
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from FleetConfig.Profile's actual components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link FleetConfig.Profile} through the TRUE canonical constructor — resolved by the
* record's own component types, not by argument count — so this never accidentally exercises a
* back-compat overload the way a literal {@code new Profile(...)} call risks doing.
*/
private static FleetConfig.Profile profileOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<FleetConfig.Profile> ctor = FleetConfig.Profile.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
@Test
void exclusionListSizeIsPinned() {
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
+ "growing exclusion list that silences failures on its own is not a guard");
}
/**
* The mutation this is built to catch: make {@code withProfile(String)}'s final constructor call
* literal at some arg count, add one more component to the record with a new back-compat
* constructor at the old arity, and the stale call silently rebinds. Every component here is real
* and non-null/non-blank, so none of it should be replaced by {@code withProfile}, except
* {@code profile} itself, which the method is documented to replace.
*/
@Test
void withProfilePreservesEveryOtherComponent() throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
FleetConfig.Profile profile = profileOf(base);
FleetConfig.Profile renamed = profile.withProfile("renamed-profile-guard");
List<String> dropped = new ArrayList<>();
int checked = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
continue;
}
checked++;
Object expected = "profile".equals(name) ? "renamed-profile-guard" : base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(renamed);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read FleetConfig.Profile." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: withProfile() was expected to carry (%s) for '%s' but returned %s — a "
+ "component silently dropped by withProfile(), the shape of the "
+ "defect this test exists to catch (its final \"return new "
+ "Profile(...)\" call binding to a back-compat constructor instead "
+ "of the true canonical one)",
name, expected, name, actual));
}
}
System.out.printf(Locale.ROOT,
"FleetConfig.Profile.withProfile() component-survival coverage — %d components, %d "
+ "checked, %d excluded, %d survived%n",
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
assertEquals(List.of(), dropped,
"withProfile() silently dropped these components: " + dropped);
}
}
@@ -30,6 +30,7 @@ import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.BiConsumer;
import java.util.function.LongSupplier;
@@ -206,6 +207,139 @@ class FleetHealthMonitorTest {
.workingSuspectAfterOrDefault());
}
// --- fleetd #386: a stall detector whose only clock freezes with a sleeping host is worse
// than a silent one — it reports "quiet" for a member that was genuinely busy for hours.
@Test void monotonicClockFrozenPastThresholdOnRealClockStillReportsStallSuspected() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
AtomicLong mono = new AtomicLong(0);
AtomicLong real = new AtomicLong(0);
FakeHerdr herdr = new FakeHerdr().withAgent("busy", "term_busy", "pane_busy", "tab_busy");
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitorWithClocks(herdr,
List.of(member("term_busy", MemberSession.State.BUSY, 0, 0)), scheduler,
mono::get, real::get, 60, 600, (_, _) -> { });
monitor.tick(); // establishes the clock baseline; nothing has diverged yet
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("state=STALL_SUSPECTED")).count());
// The host "sleeps": the monotonic clock stands completely still while the real clock
// keeps moving, past the 600s stall threshold.
real.set(TimeUnit.SECONDS.toNanos(700));
monitor.tick();
monitor.stop();
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
.contains("member=term_busy state=STALL_SUSPECTED")),
"the real clock crossed the stall threshold even though the monotonic clock never moved");
} finally {
logger.detachAppender(appender);
}
}
@Test void clockDivergenceIsLoggedOnceNotOncePerTick() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
AtomicLong mono = new AtomicLong(0);
AtomicLong real = new AtomicLong(0);
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitorWithClocks(new FakeHerdr(), List.of(), scheduler,
mono::get, real::get, 60, 600, (_, _) -> { });
monitor.tick(); // baseline: no divergence possible yet
// One sleep gap: the monotonic clock is frozen while the real clock jumps far past one
// tick interval (60s).
real.set(TimeUnit.SECONDS.toNanos(700));
monitor.tick();
// The host is awake again: both clocks advance together from here, so no more divergence.
mono.set(TimeUnit.SECONDS.toNanos(10));
real.set(TimeUnit.SECONDS.toNanos(710));
monitor.tick();
mono.set(TimeUnit.SECONDS.toNanos(20));
real.set(TimeUnit.SECONDS.toNanos(720));
monitor.tick();
monitor.stop();
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("the monotonic clock did not advance"))
.count(), "one sleep gap must produce exactly one divergence line, not one per tick");
} finally {
logger.detachAppender(appender);
}
}
/**
* fleetd #386 follow-up, added on merge. The fix carries a PER-MEMBER drift baseline, so drift
* from a sleep that happened BEFORE a member went busy is never charged to that member. The
* two tests shipped with the fix both start with the member already BUSY, so a single global
* baseline passes them — this one fails without the per-member map.
*
* <p>Order matters: the host sleeps while nothing is busy, and only then does a member take a
* turn. Its stall clock must start at zero.
*/
@Test void driftFromASleepBeforeAMemberWentBusyIsNotChargedToThatMember() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
AtomicLong mono = new AtomicLong(0);
AtomicLong real = new AtomicLong(0);
AtomicReference<List<MemberSession>> roster = new AtomicReference<>(List.of());
FakeHerdr herdr = new FakeHerdr().withAgent("busy", "term_busy", "pane_busy", "tab_busy");
var scheduler = Executors.newSingleThreadScheduledExecutor();
AgentControl agents = new AgentControl(herdr);
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, roster::get,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, mono::get, real::get, 60, 600, (_, _) -> { });
monitor.tick(); // baseline, no members yet
// The host sleeps for 700s with nobody busy: the monotonic clock stands still.
real.set(TimeUnit.SECONDS.toNanos(700));
monitor.tick();
// Awake again. Only NOW does a member start a turn, with a fresh activity stamp taken
// from the monotonic clock. Both clocks advance together from here.
mono.set(TimeUnit.SECONDS.toNanos(10));
real.set(TimeUnit.SECONDS.toNanos(710));
roster.set(List.of(member("term_busy", MemberSession.State.BUSY,
0, TimeUnit.SECONDS.toNanos(10))));
monitor.tick();
mono.set(TimeUnit.SECONDS.toNanos(20));
real.set(TimeUnit.SECONDS.toNanos(720));
monitor.tick();
monitor.stop();
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_busy state=STALL_SUSPECTED")).count(),
"the member has been busy for 10s, not 710s — the earlier sleep is not its stall");
} finally {
logger.detachAppender(appender);
}
}
private static FleetHealthMonitor monitorWithClocks(FakeHerdr herdr, List<MemberSession> roster,
java.util.concurrent.ScheduledExecutorService scheduler, LongSupplier clock,
LongSupplier realtimeClock, long intervalSeconds, long workingSuspectAfterSeconds,
BiConsumer<String, String> failTarget) {
AgentControl agents = new AgentControl(herdr);
return new FleetHealthMonitor(agents, () -> roster,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, clock, realtimeClock, intervalSeconds, workingSuspectAfterSeconds, failTarget);
}
@Test void goneMemberRecoveryLogsOnceWithoutRefiringTargetFailure() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
Level previousLevel = logger.getLevel();
@@ -357,6 +357,54 @@ class CompletionResolverTest {
"the failure carries whatever was on screen: " + waiter.getNow(null).text());
}
@Test
void aPlausibleLookingReplyInsideTheFloorStillFails() {
// fleetd#376 guard. A fix was attempted that inspected the pane inside the floor and resolved
// a COMPLETION when the text "looked like a real reply". Every cheap test for that is unsafe:
// lastAssistantBlock falls back to the WHOLE pane when there is no ⏺ marker, and a crash pane
// almost always contains sentence punctuation — in a file path, a version, or a hostname.
// This pane is the trap: it reads like a finished answer and it is a backend failure.
FakeHerdr herdr = new FakeHerdr().readText(
"Error: connection reset while loading src/main/java/Foo.java v1.2.3\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]);
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // inside the floor
resolver.resolve("term_a", turn);
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"inside the floor the verdict is always FAILED — never guess a completion from pane text");
}
@Test
void theTooFastFailureDoesNotAssertACauseItCannotKnow() {
// fleetd#376: the message used to say "most likely a backend error before any work started".
// When no error pattern matches, that cause is a guess, and a reader who believes it stops
// looking at the pane. The verdict stays FAILED; only the claim about WHY is withdrawn.
FakeHerdr herdr = new FakeHerdr().readText("I am running on opencode/mimo-v2.5-free.\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]);
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // inside the floor
resolver.resolve("term_a", turn);
String text = waiter.getNow(null).text();
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(), "still fails, still loud");
assertFalse(text.contains("most likely a backend error"),
"an unmatched fast turn must not assert a backend error: " + text);
assertTrue(text.contains("mimo-v2.5-free"), "the pane is still carried: " + text);
}
@Test
void aBusyToDoneTransitionJustOutsideTheFloorResolvesNormally() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ a real, if quick, answer\n❯ ");
@@ -1117,8 +1165,16 @@ class CompletionResolverTest {
assertTrue(waiter.isDone());
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"still a failure — the floor itself, not the pattern, is why");
assertTrue(waiter.getNow(null).text().contains("too fast to be real work"),
"a non-match inside the floor stays the generic too-fast reason: " + waiter.getNow(null).text());
// fleetd#376: this used to assert the phrase "too fast to be real work", which carried the
// claim "most likely a backend error before any work started". With no pattern matched that
// cause is a guess, so the wording was withdrawn. What this test really guards is unchanged:
// the floor alone still fails the turn, it stays generic, and it never notifies the sink.
String reason = waiter.getNow(null).text();
assertTrue(reason.contains("inside the floor"),
"a non-match inside the floor stays the generic floor reason: " + reason);
assertFalse(reason.contains("most likely a backend error"),
"a non-match must not assert a cause it did not establish: " + reason);
assertTrue(reason.contains("still starting up"), "the pane is still carried: " + reason);
assertTrue(notified.isEmpty(), "a non-match must never notify the typed sink");
}
}
@@ -108,8 +108,10 @@ class FleetMcpTest {
}
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
// fleetd #365: a resolved live send must read distinctly from a merely-queued reply —
// see replyWithNoPendingSendIsQueuedNotError below for the other case.
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "LGTM");
assertEquals("delivered", textOf(reply));
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS);
assertNotEquals(Boolean.TRUE, res.isError());
@@ -135,7 +137,7 @@ class FleetMcpTest {
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "async LGTM");
assertEquals("delivered", textOf(reply));
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
// Poll until the async send completes and reports the reply.
McpSchema.CallToolResult polled = FleetMcp.poll(messages, ticket, null);
@@ -328,9 +330,10 @@ class FleetMcpTest {
@Test
void replyWithNoPendingSendIsQueuedNotError() {
// CB-307: a reply with no open send is now queued in the inbox, not an error.
// fleetd #365: it must also no longer claim "delivered" — nothing was waiting for it.
McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", "orphan");
assertNotEquals(Boolean.TRUE, res.isError(), "a queued reply is not an error");
assertEquals("delivered", textOf(res));
assertEquals(MessageService.ReplyOutcome.QUEUED.description(), textOf(res));
// The reply is drainable by target.
var drained = messages.drainReplies("term_a");
@@ -413,7 +416,7 @@ class FleetMcpTest {
}
assertTrue(rendezvous.isWaiting("term_a"), "the answer should have reopened a waiter");
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "done");
assertEquals("delivered", textOf(reply));
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
}
@@ -8,6 +8,7 @@ import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
@@ -45,6 +46,15 @@ class EnvAllowListScrubTest {
/** Env var names appearing in command output; anything else (prompts, wrapped lines) is noise. */
private static final Pattern ENV_NAME = Pattern.compile("^([A-Za-z_][A-Za-z0-9_]*)$");
/**
* fleetd #388: the sentinel name the generated {@code .zshenv} guard uses, kept here as a
* literal rather than referencing {@link EnvAllowListScrub#SCRUB_SENTINEL} — the two tests that
* use it must still compile and run against the pre-fix production class (which has no such
* constant), so the revert-and-prove-it-fails step exercises a real assertion instead of a
* compilation error.
*/
private static final String SCRUB_SENTINEL_NAME = "_CB633_SCRUBBED";
/**
* The equality test. Expected survivors = baseline exports ∩ allowed — i.e. every survivor is
* allowed AND every allowed name that existed survives. The operator's own secret-store exports
@@ -108,6 +118,31 @@ class EnvAllowListScrubTest {
"allowed N of M with N <= M — the denominator is always reported");
}
/** A group-shared ZDOTDIR still lets the member truncate and write its pre-created receipt. */
@Test
void groupSharedScrubWritesAndReadsItsReport(@TempDir Path tmp) throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
Set<String> allowed = MemberEnvAllowList.derive(List.of());
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed, currentUserGroup());
Map<String, String> cleanParent = Map.of(
"HOME", System.getProperty("user.home"),
"PATH", "/usr/bin:/bin",
"SHELL", "/bin/zsh");
exportedNamesFromCleanParent(cleanParent, zdotdir);
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(zdotdir);
assertNotNull(report, "a group-shared completed login shell must leave a report behind");
assertTrue(report.allowed() >= 0 && report.total() >= report.allowed(),
"allowed N of M with N <= M — the denominator is always reported");
assertEquals("rw-rw----", java.nio.file.attribute.PosixFilePermissions.toString(
Files.getPosixFilePermissions(zdotdir.resolve(EnvAllowListScrub.REPORT_FILE))),
"the pre-created receipt must be group-writable");
assertEquals("rwxr-x---", java.nio.file.attribute.PosixFilePermissions.toString(
Files.getPosixFilePermissions(zdotdir)),
"group sharing must not make the ZDOTDIR directory group-writable");
}
/** Report parsing is lenient: absent file → null (no measurement), not an exception. */
@Test
void readReportReturnsNullForADirectoryWithoutOne(@TempDir Path dir) {
@@ -224,6 +259,140 @@ class EnvAllowListScrubTest {
+ "shell. A difference here means the scrub is dead on Linux.");
}
/**
* fleetd #388: the actual gap. zsh reads {@code .zshenv} always, {@code .zprofile}/
* {@code .zlogin} only for a LOGIN shell, and {@code .zshrc} only for an INTERACTIVE one — so a
* shell that is NEITHER (a bare {@code /bin/zsh} reading a script off a non-tty stdin, no
* {@code -l}, no {@code -i}) reads only {@code .zshenv} and stops. Before this fix, that shell
* never reached {@code scrub.zsh} at all: the decoy secret below would survive untouched. This
* test injects that decoy directly into the process environment (not via a sourced dotfile,
* since the whole point of the gap is that {@code .zshenv} is normally close to empty) so the
* test does not depend on any real {@code ~/.zshrc} content existing on the host.
*/
@Test
void scrubRunsInAShellThatIsNeitherLoginNorInteractive(@TempDir Path tmp) throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
Set<String> allowed = MemberEnvAllowList.derive(List.of());
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed);
Map<String, String> cleanParent = new HashMap<>(Map.of(
"HOME", System.getProperty("user.home"),
"PATH", "/usr/bin:/bin",
"SHELL", "/bin/zsh",
"USER", System.getProperty("user.name", "nobody"),
"TMPDIR", tmp.toString()));
cleanParent.put("FLEETD_TEST_DECOY_SECRET", "x"); // not on any allow-list; must be blanked
List<String> neither = List.of(); // no -l, no -i; stdin is a pipe (never a tty) either way
Set<String> baseline = exportedNamesFromCleanParent(cleanParent, null, neither);
Set<String> scrubbed = exportedNamesFromCleanParent(cleanParent, zdotdir, neither);
Set<String> expected = new TreeSet<>();
for (String name : baseline) {
if (MemberEnvAllowList.keeps(allowed, name)) {
expected.add(name);
}
}
assertTrue(baseline.contains("FLEETD_TEST_DECOY_SECRET"),
"sanity: the decoy must actually reach the un-scrubbed baseline, or this test proves "
+ "nothing");
expected.add("ZDOTDIR"); // the harness set it and it is infrastructure, so it must survive
expected.add(SCRUB_SENTINEL_NAME); // set by the new .zshenv guard once scrubbed
assertEquals(expected, scrubbed,
"a pane shell that is NEITHER login nor interactive must still be scrubbed — its "
+ "surviving exported names must EQUAL baseline ∩ allow-list, plus the "
+ "sentinel the guard sets once it has run. FLEETD_TEST_DECOY_SECRET surviving "
+ "here means the gap is still open.");
}
/**
* fleetd #388 invariants 3 and 4, which a name-set equality cannot show: a member's own tooling
* forks plain, non-login, non-interactive zsh processes for a single command (the same shape as
* the pane shell itself), and such a child must (a) keep whatever its parent deliberately set
* for it, never (b) re-run the scrub and blank it, and never (c) overwrite the pane's own
* {@code scrub-report.txt} with a description of itself instead of the pane. All three can only
* be shown by actually running a child process from within the scrubbed pane shell.
*
* <p>The pane and the child both report presence via {@code ${NAME:+present}} — empty when a
* name is unset OR blanked (exported empty), {@code present} when it is set and non-empty. No
* value is ever printed, only these two shapes and the literal word {@code set}/{@code unset}
* for the sentinel.
*/
@Test
void neitherShellChildKeepsParentVariablesAndReceiptStillDescribesThePane(@TempDir Path tmp) throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
Set<String> allowed = MemberEnvAllowList.derive(List.of());
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed);
Map<String, String> paneEnv = new HashMap<>(Map.of(
"HOME", System.getProperty("user.home"),
"PATH", "/usr/bin:/bin",
"SHELL", "/bin/zsh",
"USER", System.getProperty("user.name", "nobody"),
"TMPDIR", tmp.toString()));
paneEnv.put("FLEETD_TEST_DECOY_SECRET", "x"); // not allow-listed; the pane must blank it
paneEnv.put("ZDOTDIR", zdotdir.toAbsolutePath().toString());
// The pane's own script reports what IT sees, then forks a plain non-login, non-interactive
// child — the shape a member's own tooling uses — carrying a variable the "parent" (this
// pane) deliberately set for it, the way git sets GIT_DIR for a hook.
String outerScript = """
print -r -- "PANE_SENTINEL=${%1$s:+set}"
print -r -- "PANE_DECOY=${FLEETD_TEST_DECOY_SECRET:+present}"
FLEETD_TEST_TOOL_VAR=keep /bin/zsh <<'CHILD'
print -r -- "CHILD_LOGIN=$([[ -o login ]] && echo yes || echo no)"
print -r -- "CHILD_INTERACTIVE=$([[ -o interactive ]] && echo yes || echo no)"
print -r -- "CHILD_TOOL_VAR=${FLEETD_TEST_TOOL_VAR:+present}"
print -r -- "CHILD_DECOY=${FLEETD_TEST_DECOY_SECRET:+present}"
print -r -- "CHILD_SENTINEL=${%1$s:+set}"
CHILD
exit
""".formatted(SCRUB_SENTINEL_NAME);
ProcessBuilder pb = new ProcessBuilder("/bin/zsh"); // no -l, no -i: the pane's own shape
pb.environment().clear();
pb.environment().putAll(paneEnv);
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
Process zsh = pb.start();
zsh.getOutputStream().write(outerScript.getBytes(StandardCharsets.UTF_8));
zsh.getOutputStream().flush();
zsh.getOutputStream().close();
String stdout = new String(zsh.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
"the pane+child probe did not exit within 60s");
assertTrue(zsh.exitValue() == 0, "probe zsh exited non-zero: " + stdout);
Map<String, String> reported = new HashMap<>();
for (String line : stdout.split("\n")) {
int eq = line.indexOf('=');
if (eq > 0) {
reported.put(line.substring(0, eq).trim(), line.substring(eq + 1).trim());
}
}
assertEquals("set", reported.get("PANE_SENTINEL"),
"the pane itself is neither login nor interactive, so the .zshenv guard must have "
+ "run the scrub and exported the sentinel");
assertEquals("", reported.get("PANE_DECOY"),
"the pane must blank a non-allow-listed name — invariant 1");
assertEquals("no", reported.get("CHILD_LOGIN"), "sanity: the child must also be non-login");
assertEquals("no", reported.get("CHILD_INTERACTIVE"), "sanity: the child must also be non-interactive");
assertEquals("present", reported.get("CHILD_TOOL_VAR"),
"invariant 3: a variable the pane deliberately set for its child must survive — a "
+ "child that re-ran the scrub would have blanked it");
assertEquals("", reported.get("CHILD_DECOY"),
"a name already blanked by the pane must stay blanked in the child, never resurrected");
assertEquals("set", reported.get("CHILD_SENTINEL"),
"the child must inherit the sentinel from the pane's environment, or it would re-scrub");
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(zdotdir);
assertNotNull(report, "the pane's own scrub pass must leave a report behind");
assertTrue(report.blanked().stream().noneMatch(n -> n.startsWith("FLEETD_TEST_TOOL_VAR")),
"invariant 4: the receipt must still describe the PANE, not the child — a child that "
+ "re-ran the scrub would have rewritten this file to list its own "
+ "FLEETD_TEST_TOOL_VAR as blanked");
}
/**
* Run {@code /bin/zsh -l -i} from a clean parent and return the NAMES it has exported by prompt
* time. With {@code zdotdir} non-null, {@code ZDOTDIR} points at a generated scrub directory, so
@@ -54,6 +54,40 @@ class LeadCoordLoopTest {
assertTrue(channel.peek().isEmpty(), "and is no longer held");
}
@Test
void redeliveryOfAMessageAlreadyWrittenToThePaneIsAckedWithoutAnotherPaneWrite() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
var herdr = new FakeHerdr().agentStatus("idle");
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
loop.tick();
channel.hold(new LeadMessage("m1", PEER, SELF, "recover me"));
loop.tick();
assertEquals(1, prompts(herdr).size(), "a redelivery must not consume the lead pane twice");
assertEquals(List.of("m1", "m1"), channel.acked(), "the redelivery still needs a fresh broker ack");
}
@Test
void aRedeliveryIsAckedEvenWhileTheLeadIsMidTurn() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
var herdr = new FakeHerdr().agentStatus("idle");
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
loop.tick();
channel.hold(new LeadMessage("m1", PEER, SELF, "recover me"));
herdr.agentStatus("working");
loop.tick();
assertEquals(1, prompts(herdr).size(), "the pane is still written exactly once");
assertEquals(List.of("m1", "m1"), channel.acked(),
"a message already written to the pane must be acked even mid-turn: the mid-turn "
+ "gate exists to protect the pane, and this message needs no pane. Gating the ack "
+ "on it leaves the message held on a lead that is busy most of the time, and every "
+ "recovery redelivers it again — which is the loop this fix exists to stop");
assertTrue(channel.peek().isEmpty(), "so it is no longer held");
}
@Test
void leavesTheMessageUnackedWhenTheLeadIsMidTurn() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
@@ -92,6 +92,33 @@ class LeadMailboxTest {
}
}
@Test
void ackThrowsWhenRecoveryClearedTheHeldMessage() throws Exception {
String to = coordId("lead-recovery-ack");
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
inbox.publish(to, new LeadMessage("recovery-ack", "lead-from", to, "in flight"));
assertEquals(1, awaitPeek(inbox).size(), "the broker delivery must be held before recovery clears it");
inbox.clearHeldForRecovery();
assertThrows(IllegalStateException.class, () -> inbox.ack("recovery-ack"),
"a cleared delivery has no valid tag, so ack must report that it did not reach the broker");
}
}
@Test
void ackOfAMessageAlreadyAckedOnThisConnectionStaysQuiet() throws Exception {
String to = coordId("lead-repeat-ack");
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
inbox.publish(to, new LeadMessage("repeat-ack", "lead-from", to, "once"));
awaitPeek(inbox);
inbox.ack("repeat-ack");
inbox.ack("repeat-ack");
}
}
@Test
void duplicateMsgIdIsNotDoubleQueued() throws Exception {
String to = coordId("lead-dedup");
@@ -479,7 +479,8 @@ class MessageServiceTest {
// The worker resumes on its own (per the ask() contract) and eventually sends its real
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted");
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET, messages.reply(T, "real result"),
"the worker's real reply must still be accepted, resolving the parked async ticket");
} finally {
messages.setAskTimeoutRaceHookForTest(null);
}
@@ -767,7 +768,9 @@ class MessageServiceTest {
@Test
void replyQueuesInInboxWhenNoSendIsOpen() {
// No send is open for this session — reply should queue in the inbox.
assertTrue(messages.reply(T, "queued-text"), "reply should succeed (queued)");
// fleetd #365: this is the case that must read as QUEUED, not "delivered".
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "queued-text"),
"reply should succeed but only as queued — nothing was waiting for it");
var drained = messages.drainReplies(T);
assertEquals(1, drained.size());
@@ -780,7 +783,9 @@ class MessageServiceTest {
awaitUninterruptibly(T);
// An explicit reply resolves the open send.
assertTrue(messages.reply(T, "send-resolved"), "reply should succeed (resolved live send)");
// fleetd #365: this is the other case — RESOLVED_SEND, distinct from QUEUED above.
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, messages.reply(T, "send-resolved"),
"reply should succeed by resolving the live waiting send");
// The inbox should be empty — the reply went to the send, not the inbox.
assertTrue(messages.drainReplies(T).isEmpty(), "no reply in the inbox");
@@ -1094,8 +1099,11 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait gives up before the worker finishes resuming");
// The worker keeps working past that window and only now calls fleet_reply.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// The worker keeps working past that window and only now calls fleet_reply. The forward
// waiter answer() opened already timed out, so this resolves via the parked async ticket,
// not a live send (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1119,7 +1127,10 @@ class MessageServiceTest {
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// No live waiter (answer()'s own forward wait already timed out) — resolves the parked
// async ticket instead (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
// fleet_stop tears the worker's session down right after the reply landed — this must never
// report the misleading "the worker session was released before it replied": a reply is
@@ -1170,7 +1181,9 @@ class MessageServiceTest {
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// A live waiter is open (the forward wait above) — this resolves it directly (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND,
messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the lead's own answer() call must not throw because ask()'s timeout cleanup raced it");
@@ -1221,8 +1234,9 @@ class MessageServiceTest {
// The worker keeps working past the timeout and only now calls fleet_reply — with no live
// rendezvous waiter open (ask()'s timeout already closed it) and no new send() having
// reopened one for this target.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// reopened one for this target. So it resolves the parked async ticket (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1279,7 +1293,10 @@ class MessageServiceTest {
// real reply — reproduce that interleaving directly instead of trying to win a real race.
messages.forgetTurnForTest(turnId);
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// answer() is still waiting on its own forward waiter for the resumed turn — a live send —
// so this resolves it directly, not the async ticket (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND,
messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the primary's own answer() call must still see the worker's real reply");
@@ -1321,7 +1338,9 @@ class MessageServiceTest {
messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
try {
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// No live waiter — resolves the parked async ticket (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1353,7 +1372,8 @@ class MessageServiceTest {
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome());
assertTrue(messages.reply(T, "which task does this answer?"));
// Ambiguous — two candidates, so it must fall back to the inbox rather than guess (fleetd #365).
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "which task does this answer?"));
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(),
"an ambiguous reply must not guess ticket1");
@@ -2032,7 +2052,7 @@ class MessageServiceTest {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
assertTrue(messages.reply(T, "resolved-live"));
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, messages.reply(T, "resolved-live"));
assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded");
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
@@ -2042,14 +2062,14 @@ class MessageServiceTest {
@Test
void hasStrandedReplyIsTrueWhenNoSendWasWaiting() {
// No send is open for T — the reply queues into the inbox and is recorded as stranded.
assertTrue(messages.reply(T, "nobody was waiting"));
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "nobody was waiting"));
assertTrue(messages.hasStrandedReply(T),
"a reply with no open send strands, even though it is safely queued in the inbox");
}
@Test
void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
assertTrue(messages.reply(T, "stray"));
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "stray"));
assertTrue(messages.hasStrandedReply(T));
// The next accepted delivery for T clears the stale stranding fact — the one case the
@@ -2067,7 +2087,7 @@ class MessageServiceTest {
@Test
void hasStrandedReplyClearsOnAbandon() {
assertTrue(messages.reply(T, "stray"));
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "stray"));
assertTrue(messages.hasStrandedReply(T));
messages.abandon(T, "session released");
@@ -984,7 +984,7 @@ class ReplyPushLoopTest {
// --- metrics (CB-512) ----------------------------------------------------------------------
@Test
void successfulNudgeIncrementsDelivered() throws Exception {
void successfulNudgeIncrementsSent() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
@@ -994,11 +994,13 @@ class ReplyPushLoopTest {
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"one nudge (1 agent.prompt call) should have been sent");
// The delivered count is bumped on the scheduler thread right after the send that releases
// The sent count is bumped on the scheduler thread right after the send that releases
// the latch — settle briefly so the counter is published before we read it.
Thread.sleep(200);
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"),
"a successfully sent nudge must count as delivered");
// fleetd #365: "sent", not "delivered" — this only proves the herdr call succeeded, not
// that the primary's pane read it.
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"),
"a successfully sent nudge must count as sent");
}
@Test
@@ -1013,11 +1015,11 @@ class ReplyPushLoopTest {
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
assertEquals(0, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"));
assertEquals(0, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"));
}
@Test
void successfulTicketNudgeIncrementsDelivered() throws Exception {
void successfulTicketNudgeIncrementsSent() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
Metrics metrics = new Metrics();
@@ -1026,8 +1028,8 @@ class ReplyPushLoopTest {
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one ticket nudge should have been sent");
Thread.sleep(200);
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"),
"a successfully sent ticket nudge must count as delivered, same metric as CB-307");
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"),
"a successfully sent ticket nudge must count as sent, same metric as CB-307");
}
@Test
@@ -0,0 +1,144 @@
package dev.ltms.fleet.peer;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Fleetd #382, the same "defect factory" #357 and #358 guarded on {@code FleetConfig.withDefaults()}
* and {@link dev.ltms.fleet.session.MemberSession}'s rebuild sites — reproduced here on
* {@link SpawnRequest#withProfile(String)}.
*
* <p>{@code SpawnRequest} has back-compat constructors at arity 3 and 5 alongside its canonical
* arity-6 constructor. Before this ticket, {@code CompositePeerLauncher} routed a profile by
* building a fresh {@code SpawnRequest} from a literal {@code new SpawnRequest(...)} call listing
* six of the original request's own accessors. That call is only correct because it happens to
* name exactly six arguments today — add a 7th component and the established back-compat pattern
* (a new constructor at the old, now-shorter arity) and a call one argument short of the new
* canonical arity would silently rebind to that back-compat constructor, dropping the new
* component on every profile-routed spawn without any compile error. {@link #withProfile} replaces
* that literal call, so this test guards the ONE rebuild site instead of a call site scattered
* through a launcher.
*
* <p>The check below builds a {@link SpawnRequest} through the TRUE canonical constructor —
* resolved by the record's own component types via {@code getDeclaredConstructor}, never by
* argument count, so it can never itself land on a back-compat overload — with a real, distinctive,
* non-null value in every component, calls {@link SpawnRequest#withProfile(String)}, and asserts
* every component the method is not documented to change survives unchanged, while {@code
* profileName} comes back as the new value it was given. A component that comes back anything else
* was silently dropped — the shape of the defect this test exists to catch.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
* {@link #exclusionListSizeIsPinned()}, for the same reason the other two guards pin theirs at
* zero: a checker whose escape hatch can grow to silence a failure is not a checker. Every one of
* {@link SpawnRequest}'s 6 current components has a real, non-null value here and none is excluded.
*/
class SpawnRequestWithProfilePreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = SpawnRequest.class.getRecordComponents();
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null value per component — none of the 6 is excluded. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profileName", "profile-guard");
v.put("requestedCwd", "/wt/requested-guard");
v.put("callerCwd", "/wt/caller-guard");
v.put("sessionName", "session-guard");
v.put("resumeSessionId", "resume-guard");
v.put("role", MemberRole.REVIEWER);
assertNamesMatchComponents(v);
return v;
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
* to add a new component here fails this assertion by name, rather than silently checking one
* component fewer than the record has.
*/
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from SpawnRequest's actual components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link SpawnRequest} through the TRUE canonical constructor — resolved by the
* record's own component types, not by argument count — so this never accidentally exercises a
* back-compat overload the way a literal {@code new SpawnRequest(...)} call risks doing.
*/
private static SpawnRequest requestOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<SpawnRequest> ctor = SpawnRequest.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
@Test
void exclusionListSizeIsPinned() {
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
+ "growing exclusion list that silences failures on its own is not a guard");
}
@Test
void withProfilePreservesEveryOtherComponent() throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
SpawnRequest request = requestOf(base);
SpawnRequest result = request.withProfile("profile-updated");
Map<String, Object> expectedOverrides = Map.of("profileName", "profile-updated");
List<String> dropped = new ArrayList<>();
int checked = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
continue;
}
checked++;
Object expected = expectedOverrides.containsKey(name) ? expectedOverrides.get(name) : base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(result);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read SpawnRequest." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: withProfile() was expected to carry (%s) for '%s' but returned %s — a "
+ "component silently dropped by withProfile(), the shape of the defect "
+ "this test exists to catch (its final \"return new SpawnRequest(...)\" "
+ "call binding to a back-compat constructor instead of the true "
+ "canonical one)",
name, expected, name, actual));
}
}
System.out.printf(Locale.ROOT,
"SpawnRequest.withProfile() component-survival coverage — %d components, %d checked, "
+ "%d excluded, %d survived%n",
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
assertEquals(List.of(), dropped,
"withProfile() silently dropped these components: " + dropped);
}
}
@@ -476,6 +476,11 @@ class FleetAppTest {
Thread.sleep(200);
HttpResponse<String> reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}");
assertEquals(200, reply.statusCode());
// fleetd #365: "delivered" used to be unconditionally true; a live send was actually waiting
// here, so this is the case where it must genuinely read true, with outcome naming why.
JsonNode replyBody = mapper.readTree(reply.body());
assertEquals(true, replyBody.get("delivered").asBoolean());
assertEquals("resolved_send", replyBody.get("outcome").asText());
HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS);
assertEquals(200, res.statusCode());
@@ -527,6 +532,11 @@ class FleetAppTest {
int port = startHealthy();
HttpResponse<String> res = postJson(port, "/sessions/term_a/reply", "{\"content\":\"orphan\"}");
assertEquals(200, res.statusCode());
// fleetd #365: nothing was waiting, so "delivered" must now read false, not the old
// unconditional true — outcome names this as queued.
JsonNode resBody = mapper.readTree(res.body());
assertEquals(false, resBody.get("delivered").asBoolean());
assertEquals("queued", resBody.get("outcome").asText());
// The queued reply is drainable.
HttpResponse<String> drain = req(port, "GET", "/sessions/term_a/replies");
@@ -1532,26 +1532,52 @@ class GitWorktreesTest {
* for repo setup.
*
* <p>Scope, measured on the fleetd #369 merge and narrower than an earlier version of this
* comment claimed: this protects the 5 {@link #seedingGitWorktrees} sites plus — through
* comment claimed: this protects the {@link #seedingGitWorktrees} call sites plus — through
* {@link #gitProcessBuilder} — every {@code git} subprocess the TEST itself starts. It does
* NOT cover the other 53 {@code new GitWorktrees(...)} constructions in this file, which pass
* NOT cover the {@code new GitWorktrees(...)} constructions elsewhere in this file that pass
* no env override, so a production instance built that way still inherits the JVM's real
* environment. Stripping this override from {@code seedingGitWorktrees} leaves the class green
* both with and without the poison command above, so that half is currently unpinned.
* environment. (Re-measured for fleetd #373, on this file as it stands here: 4 call sites go
* through {@link #seedingGitWorktrees(Path, String, Path)} — not 5, an earlier count this
* comment and fleetd #373's own ticket text both repeated without re-running it — out of 59
* total {@code new GitWorktrees(...)} occurrences, one of which is the shared construction
* inside {@link #seedingGitWorktrees(Path, String, Map)} itself. This class-wide count moves
* every time a test is added, so treat any number here as a snapshot, not a fact to cite
* without recounting.) Stripping the {@code gitEnv} override from a {@link
* #seedingGitWorktrees} call site leaves the class green both with and without the poison
* command above for that call site's OWN test, so that half was unpinned until fleetd #373
* added {@link #seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory}
* below, which asserts the property directly instead of relying on a poisoned real machine.
*/
private static Map<String, String> hermeticGitEnv(Path tmp) {
return hermeticGitEnvAt(tmp.resolve("hermetic-xdg-config-home-" + System.nanoTime()));
}
/** Same isolation as {@link #hermeticGitEnv(Path)}, with an explicit {@code XDG_CONFIG_HOME}
* instead of a fresh nanoTime-unique one under {@code tmp} — used by
* {@link #seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory}
* (fleetd #373) so it can pre-populate that directory with a marker BEFORE the production
* {@link GitWorktrees} instance reads it, something the random per-call name from
* {@link #hermeticGitEnv(Path)} makes impossible to predict from outside. */
private static Map<String, String> hermeticGitEnvAt(Path xdgConfigHome) {
return Map.of(
"GIT_CONFIG_GLOBAL", "/dev/null",
"GIT_CONFIG_SYSTEM", "/dev/null",
"GIT_TERMINAL_PROMPT", "0",
"XDG_CONFIG_HOME", tmp.resolve("hermetic-xdg-config-home-" + System.nanoTime()).toString());
"XDG_CONFIG_HOME", xdgConfigHome.toString());
}
/** {@link GitWorktrees}'s full test seam, with a {@code memberSkillsSource} and no other
* overrides — the shape every seeding test below needs, isolated via {@link #hermeticGitEnv}. */
private static GitWorktrees seedingGitWorktrees(Path root, String memberSkillsSource, Path tmp) {
return new GitWorktrees(root.toString(), null, _ -> {}, null, null, memberSkillsSource,
hermeticGitEnv(tmp));
return seedingGitWorktrees(root, memberSkillsSource, hermeticGitEnv(tmp));
}
/** Same shape as {@link #seedingGitWorktrees(Path, String, Path)}, taking an already-built
* {@code gitEnv} directly rather than computing one via {@link #hermeticGitEnv(Path)} — lets
* fleetd #373's test drive the exact production construction a real member spawn uses, with a
* {@code gitEnv} it has already pre-populated a marker into. */
private static GitWorktrees seedingGitWorktrees(Path root, String memberSkillsSource, Map<String, String> gitEnv) {
return new GitWorktrees(root.toString(), null, _ -> {}, null, null, memberSkillsSource, gitEnv);
}
/** Acceptance criterion 2 (part 1): a worktree with no {@code .claude/} at all gets the skill
@@ -1784,6 +1810,63 @@ class GitWorktreesTest {
+ "after skill seeding ran — got:\n" + porcelain);
}
/**
* fleetd #373. Pins the production seam that fleetd #362 review finding 2 protects: {@link
* GitWorktrees#previouslyEffectiveExcludesFileContent}'s XDG-fallback branch reads {@code
* XDG_CONFIG_HOME}/{@code HOME} straight in Java, not through a {@code git} subprocess, so
* {@code gitEnv} — the constructor seam every {@link #seedingGitWorktrees} instance in this
* class is built with — is the ONLY thing that can isolate it. A mutation run during the
* fleetd #372/#369 merge found this unpinned: replacing {@code hermeticGitEnv(tmp)} with
* {@code null} in {@link #seedingGitWorktrees(Path, String, Path)} left every test in this
* class green — the 56 tests that would fail against a real machine's poisoned {@code
* XDG_CONFIG_HOME} were fixed by fleetd #369's subprocess-level isolation, but none of them
* looks at what THIS Java-side read resolves, so deleting the override stays invisible.
*
* <p>This test asserts the PROPERTY, not the constructor argument: a {@link GitWorktrees}
* built for seeding — through the very same {@link #seedingGitWorktrees(Path, String, Map)}
* construction every other seeding test in this class goes through — must resolve the
* excludes-file fallback inside its own throwaway {@code gitEnv}-supplied directory. It needs
* NO externally-set poisoned environment variable: the marker pattern below is written ONLY
* inside a throwaway {@code XDG_CONFIG_HOME} this test controls directly (bypassing {@link
* #hermeticGitEnv(Path)}'s unpredictable nanoTime-named directory, via {@link
* #hermeticGitEnvAt}, so the marker can be in place before the production instance ever reads
* it), reachable ONLY through the {@code gitEnv} seam. If that seam is stripped, the
* production code instead falls back to resolving the REAL {@code XDG_CONFIG_HOME}/{@code
* HOME} of the machine running the test — which does not carry this marker — so the marker
* file below shows up as untracked and the assertion fails on any machine, with no poison
* command required. See the PR body for the pasted failure from actually running that
* mutation (removing the {@code gitEnv} override from this test's own construction).
*/
@Test
void seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory(@TempDir Path tmp)
throws Exception {
Path xdgConfigHome = tmp.resolve("cb373-xdg-config-home");
Files.createDirectories(xdgConfigHome.resolve("git"));
Files.writeString(xdgConfigHome.resolve("git").resolve("ignore"), "cb373-xdg-fallback-marker\n");
Map<String, String> gitEnv = hermeticGitEnvAt(xdgConfigHome);
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
GitWorktrees seeding = seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), gitEnv);
String wt = seeding.add(repo.toString(), "cb-373-xdg-seam", "HEAD");
assertEquals("IMPLEMENTER SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
"fixture check — the skill really was seeded, so previouslyEffectiveExcludesFileContent ran");
Files.writeString(Path.of(wt, "cb373-xdg-fallback-marker"),
"would only be invisible to git status if the fallback resolved THIS throwaway "
+ "XDG_CONFIG_HOME rather than the real machine's\n");
String porcelain = fullStatus(Path.of(wt));
assertEquals("", porcelain,
"the marker pattern lives only in this test's throwaway XDG_CONFIG_HOME; git "
+ "status must still be empty, proving the production seam resolved the "
+ "excludes-file fallback through the gitEnv seam rather than the JVM's "
+ "real environment — got:\n" + porcelain);
}
/**
* fleetd #369, acceptance criterion 4 — make the fix hard to undo by accident. Every git
* subprocess this class starts is required to go through {@link #gitProcessBuilder}, the one
@@ -0,0 +1,190 @@
package dev.ltms.fleet.session;
import dev.ltms.fleet.peer.CharterReceipt;
import dev.ltms.fleet.peer.MemberRole;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import java.util.function.Function;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Fleetd #358, the same "defect factory" #357 guarded on {@code FleetConfig.withDefaults()}
* (see {@code FleetConfigWithDefaultsPreservesEveryComponentTest}), reproduced here on
* {@link MemberSession} — the worse of the two sibling cases named in #358, because this record
* has FIVE independent rebuild sites instead of one: {@link MemberSession#withState},
* {@link MemberSession#withActivity}, {@link MemberSession#bumpTurn},
* {@link MemberSession#withAgentSessionId} and {@link MemberSession#withFailureReason} each end in
* their own literal {@code new MemberSession(...)} call. Add a 16th component and add the
* established back-compat constructor at the old (15-arg) arity, and every one of those five
* literal calls becomes a legal match for that new overload — silently dropping the new component,
* independently, on whichever of the five paths a missed update leaves behind. That is harder to
* spot than #357's single call site: the field would survive through some transitions and vanish
* through others.
*
* <p>Each check below builds one {@link MemberSession} through the TRUE canonical constructor —
* resolved by the record's own component types via {@code getDeclaredConstructor}, never by
* argument count, so it can never itself land on a back-compat overload — with a real, distinctive,
* non-null value in every component, calls the real rebuild method under test, and asserts every
* component the method is not documented to change survives unchanged, while the component(s) it IS
* documented to change come back as the new value it was given. A component that comes back
* anything else was silently dropped or lost — the shape of the defect this test exists to catch.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
* {@link #exclusionListSizeIsPinned()}, for the same reason {@code FleetConfig}'s guard pins its own
* exclusion list at zero: a checker whose escape hatch can grow to silence a failure is not a
* checker. Every one of {@link MemberSession}'s 15 current components has a real, non-null,
* non-blank value here and none is excluded.
*/
class MemberSessionRebuildPreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = MemberSession.class.getRecordComponents();
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null value per component — none of the 15 is excluded. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("paneId", "pane-guard");
v.put("terminalId", "term-guard");
v.put("profile", "profile-guard");
v.put("role", MemberRole.REVIEWER);
v.put("cwd", "/wt/guard");
v.put("ownerTerminal", "owner-guard");
v.put("spawnedAtNanos", 111_111L);
v.put("lastActivityAtNanos", 222_222L);
v.put("turnCount", 7);
v.put("state", MemberSession.State.BUSY);
v.put("worktree", "/wt/guard-tree");
v.put("branch", "worker/guard-branch");
v.put("charterReceipt", new CharterReceipt(
MemberRole.DEV, "profile-guard", "fleet.charters.dev", "deadbeefguard", 42));
v.put("agentSessionId", "agent-guard");
v.put("failureReason", "reason-guard");
assertNamesMatchComponents(v);
return v;
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
* to add a new component here fails this assertion by name, rather than silently checking one
* component fewer than the record has.
*/
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from MemberSession's actual components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link MemberSession} through the TRUE canonical constructor — resolved by the
* record's own component types, not by argument count — so this never accidentally exercises a
* back-compat overload the way a literal {@code new MemberSession(...)} call risks doing.
*/
private static MemberSession sessionOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<MemberSession> ctor = MemberSession.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
@Test
void exclusionListSizeIsPinned() {
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
+ "growing exclusion list that silences failures on its own is not a guard");
}
/**
* Shared check for one rebuild site: build a base session with a real value in every component,
* call {@code rebuild}, and assert every component comes back equal to {@code expectedOverrides}
* when named there, or equal to the base value otherwise. Prints the same denominator style as
* {@code FleetConfigWithDefaultsPreservesEveryComponentTest}.
*/
private void checkRebuildSite(String siteName, Function<MemberSession, MemberSession> rebuild,
Map<String, Object> expectedOverrides) throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
MemberSession session = sessionOf(base);
MemberSession result = rebuild.apply(session);
List<String> dropped = new ArrayList<>();
int checked = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
continue;
}
checked++;
Object expected = expectedOverrides.containsKey(name) ? expectedOverrides.get(name) : base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(result);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read MemberSession." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: %s() was expected to carry (%s) for '%s' but returned %s — a component "
+ "silently dropped by %s(), the shape of the defect this test exists "
+ "to catch (its final \"return new MemberSession(...)\" call binding "
+ "to a back-compat constructor instead of the true canonical one)",
name, siteName, expected, name, actual, siteName));
}
}
System.out.printf(Locale.ROOT,
"MemberSession.%s() component-survival coverage — %d components, %d checked, %d "
+ "excluded, %d survived%n",
siteName, COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
checked - dropped.size());
assertEquals(List.of(), dropped,
siteName + "() silently dropped these components: " + dropped);
}
@Test
void withStatePreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("withState", s -> s.withState(MemberSession.State.FAILED),
Map.of("state", MemberSession.State.FAILED));
}
@Test
void withActivityPreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("withActivity", s -> s.withActivity(999_999L),
Map.of("lastActivityAtNanos", 999_999L));
}
@Test
void bumpTurnPreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("bumpTurn", s -> s.bumpTurn(999_999L),
Map.of("lastActivityAtNanos", 999_999L, "turnCount", 8));
}
@Test
void withAgentSessionIdPreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("withAgentSessionId", s -> s.withAgentSessionId("agent-updated"),
Map.of("agentSessionId", "agent-updated"));
}
@Test
void withFailureReasonPreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("withFailureReason", s -> s.withFailureReason("reason-updated"),
Map.of("failureReason", "reason-updated"));
}
}