Compare commits

..

6 Commits

Author SHA1 Message Date
Dai Ha 464dbc0930 fleetd#341: a per-name guard so a later spawn's different unprotected name still warns
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Successful in 1m29s
unprotectedGapLogged was one AtomicBoolean guarding two WARN branches in
logCredentialGap that name different env var names (the allow-list
keptByDerivedList branch, and warnGapUnprotected's deny-by-default /
non-zsh-fallback branch). memberCredentials is a live, re-read-per-spawn
supplier, so between two spawns a policy reload can change which names are
in the gap: spawn 1 warns about name A and trips the shared flag, and
spawn 2's gap containing a different name B never gets its WARN.

Replace the AtomicBoolean with unprotectedGapNamesWarned, a
ConcurrentHashMap-backed Set<String> guard keyed per name (same shape as
OpenCodeLauncher.modelCheckSkippedWarned), so each distinct credential-shaped
name is warned about exactly once, ever, regardless of which branch or
which spawn first reports it. allowListGapLogged (the separate INFO guard,
#192) is untouched. Neither WARN's wording changed.
2026-09-04 15:59:43 +07:00
Dai Ha eee4d576a2 #333: the Cold doc bullet listed four keys, COLD_KEYS has five
CI / build (push) Successful in 1m52s
CI / contract (push) Successful in 2m8s
memberHerdrSocket was missing from the prose. The #333 worker spotted it and
correctly left it alone as outside its scope.

Fixed by pointing the bullet at COLD_KEYS instead of re-listing its contents,
so the prose and the set cannot drift apart a second time.
2026-09-04 15:43:39 +07:00
Dai Ha b4f9d7f53a Merge #333: fleet: is a split key, and split membership now proves a reporting branch exists 2026-09-04 15:39:02 +07:00
Dai Ha 4aa1fae296 #329: a null task in answer() is not only "never an async ticket"
CI / contract (push) Successful in 59s
CI / build (push) Successful in 2m26s
The comment that landed with #329 said a null task means the turn was never
an async ticket. That is wrong, and it makes the guard read as complete.

A genuine async ticket also reaches answer() with task == null. ask() runs
clearAsyncQuestion(turnId, true) in its catch block, which drops the
asyncTasksByTurn entry, while rendezvous.closeAsk(turnId) runs later, in its
finally. Between the two the ask is still answerable and the map entry is
already gone, so answer()'s lookup returns null and the ticket is stranded.

Measured with a throwaway probe firing only that first half: answer() reported
REPLIED while the ticket stayed PENDING with a null reply. The probe used
forgetTurnForTest, so it omits markAskTimedOut; that cannot change the outcome,
because askTimedOut is read only by askAnsweredAsyncTasks, which reply() never
reaches while answer()'s own waiter is live.

Open as fleetd #334. The comment now says so.
2026-09-04 15:33:57 +07:00
Dai Ha 0c865032f9 Merge #329: log an exception thrown after a ticket resolves, complete the ticket from the task answer() already holds, and read orphan.turnId once 2026-09-04 15:26:45 +07:00
Dai Ha ea41bbf6b9 fleetd#329: fix silent async-ticket bugs in MessageService (F1/F2/F3)
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m28s
F2 (sendAsync executor catch): log when completeExceptionally returns
false, so an exception thrown after finishAsyncTask already completed
the ticket's future is no longer silently lost.

F1 (answer()'s stranded async ticket): reuse the Task reference answer()
already looked up before rendezvous.answerAsk(), instead of a second
asyncTasksByTurn lookup by turnId in finishAsyncTask. The second lookup
raced ask()'s unlocked timeout cleanup, which could forget turnId first
and leave the ticket stuck PENDING even though answer() itself returned
REPLIED. The #282 chained-ask guard is unaffected: it is still keyed on
result.outcome() == QUESTION, not on this lookup. Removed the now-unused
finishAsyncTask(String, Reply) overload.

F3 (reply()'s orphan recovery path): read orphan.turnId once instead of
twice, closing the same double-read shape fleetd #324 fixed in
finishAsyncTask.

Each fix has its own test plus a test-only race hook (mirroring #324's
finishAsyncTaskRaceHook) to force the exact interleaving deterministically.
Mutation-tested each fix by reverting it, confirming the real failure
(swallowed exception / PENDING ticket / NullPointerException), then
restoring it.

mvn clean install: Tests run: 1345, Failures: 0, Errors: 0, Skipped: 0,
BUILD SUCCESS.
2026-09-04 15:19:38 +07:00
5 changed files with 403 additions and 34 deletions
@@ -121,10 +121,12 @@ import java.util.function.Supplier;
* reload would leave the operator worse off than today. {@link Outcome#split()} names the
* key and says which half is which each time, rather than trying to score "how changed" a
* mixed key is or handle "both halves changed in one reload" as a special case.</li>
* <li><strong>Cold</strong> — cannot change at all under a running daemon: {@code bind:},
* {@code herdrSocket:}, {@code broker:} and {@code auth:}. The socket is bound, the broker
* <li><strong>Cold</strong> — cannot change at all under a running daemon. All five of
* {@link #COLD_KEYS}: {@code bind:}, {@code herdrSocket:}, {@code memberHerdrSocket:},
* {@code broker:} and {@code auth:}. The sockets are already connected, the broker
* connection is open, and the auth mode decides who may reach the port that is already
* listening.</li>
* listening. This bullet omitted {@code memberHerdrSocket:} until fleetd #333 — say "all
* five of COLD_KEYS" rather than re-listing them, so prose and set cannot drift again.</li>
* </ul>
*
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333).</strong>
@@ -1665,30 +1665,59 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/**
* Guards the {@code effectiveAllowed == null} branch of {@link #logCredentialGap} — the
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — to one
* WARN per launcher instance, not one per spawn.
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — AND
* the {@code effectiveAllowed != null} / {@code keptByDerivedList} branch, the allow-list case
* where a name is in the gap but the derived allow-list keeps it anyway. Both branches log the
* same severity (WARN) about the same fact — a name genuinely reaching a member pane
* unprotected — so they share this one guard, keyed per NAME rather than per launcher instance:
* each credential-shaped name that is ever reported unprotected gets exactly one WARN, however
* many spawns see it and whichever of the two branches first reports it.
*
* <p>fleetd #341: {@code memberCredentials} is a live, re-read-per-spawn supplier, so the
* policy — and so the gap's actual member names — can change between two spawns on the same
* launcher. Before this fix the guard was a single {@code AtomicBoolean} tripped by either
* branch: spawn 1 could warn about name A and trip the flag, and a later spawn's gap containing
* a different name B would never be reported, even though B is just as unprotected as A was.
* {@code AtomicBoolean} could not express "once per distinct name" at all — only "once, ever,
* for whichever name got there first" — so this is a {@code Set<String>} guard instead, the
* same shape {@link OpenCodeLauncher#modelCheckSkippedWarned} already uses for its own
* once-per-distinct-thing WARN. {@link #add}'s return value (true only the first time a name is
* added) is what turns "log the whole gap" into "log only the names never warned about before".
*
* <p>Bounded by construction: every name added here first passed {@link
* #CREDENTIAL_SHAPED_NAME}'s filter over {@link #hostEnvNames}, i.e. it is an actual
* environment variable name from the daemon's own process — a small, OS-bounded set (the host
* environment has, in practice, tens to a few hundred entries), not an attacker- or
* request-controlled input. So this set cannot grow past "however many distinct credential-
* shaped names this host's environment has ever held across this launcher's lifetime," which is
* effectively fixed for the life of one daemon process — no separate cap is needed.
*
* <p>CB-633 follow-up (#192): kept SEPARATE from {@link #allowListGapLogged} on purpose.
* {@code memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change
* between two spawns on the same launcher. A single shared flag would let a harmless allow-list
* INFO on spawn 1 permanently suppress the real deny-by-default WARN a later spawn deserves —
* the report that matters most getting hidden by the report that doesn't. Two flags mean each
* report kind fires exactly once, independent of what the other kind already logged.
* report kind (WARN vs. INFO) fires independently of what the other kind already logged; within
* the WARN kind itself, the set above further separates by name, for the same reason.
*/
private final AtomicBoolean unprotectedGapLogged = new AtomicBoolean();
private final Set<String> unprotectedGapNamesWarned = ConcurrentHashMap.newKeySet();
/**
* Guards the {@code effectiveAllowed != null} branch of {@link #logCredentialGap} — the
* allow-list-scrub-covered report — to one INFO per launcher instance. See {@link
* #unprotectedGapLogged}'s javadoc for why this is a separate flag rather than a shared one.
* #unprotectedGapNamesWarned}'s javadoc for why this is a separate flag rather than a shared
* one; unlike that guard it stays a per-instance {@code AtomicBoolean}, not a per-name set —
* fleetd #341 fixed the WARN-vs-WARN suppression, not this INFO's own one-shot shape, which was
* not reported as broken and is out of that ticket's scope.
*/
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
/**
* fleetd #185 stage 2: guards {@link #warnUnknownMemberEnvironment} to one WARN per launcher
* instance, not one per spawn — the same one-per-instance shape as {@link #unprotectedGapLogged}
* and {@link #allowListGapLogged}, kept as its own flag for the same reason those two are split:
* this mode is orthogonal to which of the other two branches would otherwise have fired.
* instance, not one per spawn — the same one-shot shape {@link #unprotectedGapNamesWarned} and
* {@link #allowListGapLogged} guard their own branches with, kept as its own flag for the same
* reason those two are split: this mode is orthogonal to which of the other two branches would
* otherwise have fired.
*/
private final AtomicBoolean unknownMemberEnvironmentWarned = new AtomicBoolean();
@@ -1816,14 +1845,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
List<String> blankedByScrub = gap.stream()
.filter(name -> !MemberEnvAllowList.keeps(effectiveAllowed, name))
.toList();
if (!keptByDerivedList.isEmpty() && unprotectedGapLogged.compareAndSet(false, true)) {
// fleetd #341: filter to names this guard has never warned about before — not just
// "isEmpty" on the whole branch — so a name this spawn's gap shares with an EARLIER
// spawn's (already-warned) gap does not re-print, while a name unique to THIS gap still
// does, whichever of the two WARN branches reported it first.
List<String> newlyUnprotected = keptByDerivedList.stream()
.filter(unprotectedGapNamesWarned::add)
.toList();
if (!newlyUnprotected.isEmpty()) {
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
+ "known: nor allow: — the derived allow-list keeps them anyway (a profile's "
+ "gitTokenEnv/gitHostEnv/tokenEnv/env: names one, or this spawn injects it), "
+ "so every member pane inherits them UNBLOCKED — {}. Add each to "
+ "memberCredentials.known (or .allow if a member legitimately needs it), or "
+ "remove it from whatever profile setting derives it in.",
keptByDerivedList.size(), keptByDerivedList);
newlyUnprotected.size(), newlyUnprotected);
}
if (!blankedByScrub.isEmpty() && allowListGapLogged.compareAndSet(false, true)) {
log.info("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
@@ -1836,12 +1872,18 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/** The deny-by-default (and allow-list non-zsh fallback) WARN — unchanged byte-for-byte by #192. */
private void warnGapUnprotected(List<String> gap) {
if (unprotectedGapLogged.compareAndSet(false, true)) {
// fleetd #341: same "only the names never warned before" filter as the sibling branch in
// logCredentialGap above — see unprotectedGapNamesWarned's javadoc. Both branches share
// this one guard because both report the exact same fact (a name reaching a member pane
// unprotected) at the exact same severity; keying it by name is what lets a later spawn's
// DIFFERENT name still get its own WARN after an earlier spawn's already fired.
List<String> newlyUnprotected = gap.stream().filter(unprotectedGapNamesWarned::add).toList();
if (!newlyUnprotected.isEmpty()) {
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
+ "known: nor allow: — every member pane inherits them UNBLOCKED — {}. "
+ "Add each to memberCredentials.known (blocked by default) or .allow "
+ "(if a member legitimately needs it).",
gap.size(), gap);
newlyUnprotected.size(), newlyUnprotected);
}
}
@@ -466,8 +466,18 @@ public final class MessageService {
if (candidates.size() == 1) {
Task orphan = candidates.get(0);
if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
if (orphan.turnId != null) {
asyncTasksByTurn.remove(orphan.turnId, orphan);
// fleetd #329 (F3): read orphan.turnId once. It used to be read twice, under no
// lock — once for this null check, once as the removal key — the same double-read
// shape fleetd #324 fixed in finishAsyncTask. clearAsyncQuestion's unlocked
// forgetTurn=true path (ask()'s timeout cleanup) can null the field between the two
// reads; capturing it once removes the torn read here too.
String turnId = orphan.turnId;
if (turnId != null) {
if (replyOrphanTurnIdRaceHookForTest != null) {
// Test-only (fleetd #329, F3): see the field's own javadoc.
replyOrphanTurnIdRaceHookForTest.run();
}
asyncTasksByTurn.remove(turnId, orphan);
}
count(FleetMetrics.REPLIES, "path", "async-recovered");
return true; // the ticket itself took it — no inbox stranding at all
@@ -1002,13 +1012,49 @@ public final class MessageService {
// Measured when #282 was merged: this guard is DEFENCE IN DEPTH, not the thing
// that makes the chained ask work. ask() calls markAsyncQuestion (:860) before
// resolveQuestion (:861), so by the time this thread wakes, the task has already
// moved to the new turnId and finishAsyncTask(oldTurnId, ...) finds nothing. Removing
// this guard alone leaves the test green. Keep it anyway: it mirrors sendAsync's
// sibling guard, and that sibling's own comment (:1017) warns the two orderings are
// not something to rely on. Do NOT delete it as dead code without re-checking that
// ordering, and do not treat it as the sole protection either.
if (result.outcome() != Outcome.QUESTION) {
finishAsyncTask(turnId, result);
// moved to the new turnId — this outcome check, not a Task lookup, is what tells the
// two cases apart (see fleetd #329 below). Removing this guard alone leaves the test
// green. Keep it anyway: it mirrors sendAsync's sibling guard, and that sibling's own
// comment (:1017) warns the two orderings are not something to rely on. Do NOT delete
// it as dead code without re-checking that ordering, and do not treat it as the sole
// protection either.
//
// fleetd #329 (F1): complete the SAME Task object this method already looked up at
// :991, rather than re-resolving it from turnId a second time. The old
// finishAsyncTask(turnId, result) did its own asyncTasksByTurn.get(turnId) here, and
// that second lookup races ask()'s own timeout path: ask()'s ticket.answer().get(...)
// can time out (or lose that exact race) at essentially the same instant this method's
// rendezvous.answerAsk(turnId, ...) above already succeeded, running
// markAskTimedOut + clearAsyncQuestion(turnId, true) with no lock at all — which
// forgets turnId (removes it from asyncTasksByTurn, nulls Task.turnId) before this
// thread ever gets here. The worker's real reply then arrived, this method's own wait
// woke up with it, and the by-turnId lookup found nothing: the async ticket's future
// was never completed, so fleet_poll{ticket} stayed PENDING forever even though
// answer() itself correctly returned REPLIED. Reusing the reference captured at :991
// — before that race window opens — sidesteps the second lookup entirely: it is the
// identical Task whatever asyncTasksByTurn or Task.turnId say by the time we reach
// this point, so finishAsyncTask(Task, Reply) can still complete its future and detach
// it using whatever turnId it now reads. The #282 chained-ask case is unaffected
// because it is still gated purely by result.outcome() == QUESTION above, which does
// not depend on this lookup — widening what "task" means here cannot complete a ticket
// the chained ask deliberately left open.
//
// A null task is NOT only "this was never an async ticket". That reading was in this
// comment when #329 merged and it is wrong. A genuine async ticket also lands here
// with task == null, because ask()'s timeout path runs clearAsyncQuestion(turnId,
// true) — which drops the asyncTasksByTurn entry — in its catch block, while
// rendezvous.closeAsk(turnId) runs later, in its finally. Between those two the ask
// is still answerable but the map entry is already gone, so the lookup at :991
// returns null and this ticket is never completed. Measured on 2026-09-04: a probe
// firing only that first half before answer() runs printed
// "answer=REPLIED phase=PENDING reply=null" — the same stranded ticket #329 set out
// to fix, one step earlier in the same race. The probe used forgetTurnForTest, which
// omits ask()'s markAskTimedOut; that cannot change the outcome, because askTimedOut
// is read only by askAnsweredAsyncTasks, and reply() never reaches it while this
// method's own waiter is live. So #329 narrows this window rather than closing it.
// Open as fleetd #334 — do not read this guard as complete.
if (result.outcome() != Outcome.QUESTION && task != null) {
finishAsyncTask(task, result);
}
return result;
} catch (TimeoutException e) {
@@ -1059,7 +1105,7 @@ public final class MessageService {
// this fires exactly once, from whichever path completes it: finishAsyncTask(task, result)
// below on any non-QUESTION outcome of send() — a worker's fleet_reply, the CB-106
// completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same
// finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION
// finishAsyncTask reached via answer()'s own finishAsyncTask(task, result) once a QUESTION
// is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516
// abandon() on teardown. Without this, MessageService.reply's rendezvous fast path (the
// one an async ticket always takes) never told the push loop anything happened — see the
@@ -1079,7 +1125,22 @@ public final class MessageService {
finishAsyncTask(task, result);
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
// fleetd #329 (F2): finishAsyncTask above already completes task.future — on its very
// first line — before doing anything else, so anything that throws afterward (inside
// finishAsyncTask's own cleanup, or from a future addition to this try block) lands
// here with the future already resolved. completeExceptionally on an already-completed
// future is a silent no-op: it returns false and does nothing, so without the check
// below the exception simply vanished — no log, no metric, nothing. Measured (see the
// ticket): temporarily reintroducing the fleetd #324 NPE reproduced 19 real exceptions
// on the ordinary path, with 74/74 tests staying green and not one log line produced.
// Do not stop completing the future first — finishAsyncTask completing before it
// cleans up is what makes a late failure harmless to the ticket's own result — only
// add the missing visibility for the case where that step, or whatever ran after it,
// has already lost the race to report through the future.
if (!task.future.completeExceptionally(t)) {
log.error("async send {} -> {} threw after its ticket was already resolved",
ticket, target, t);
}
}
});
pruneTerminalTickets();
@@ -1247,6 +1308,10 @@ public final class MessageService {
*/
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
if (afterFinishAsyncTaskCompleteHookForTest != null) {
// Test-only (fleetd #329, F2): see the field's own javadoc.
afterFinishAsyncTaskCompleteHookForTest.run();
}
String turnId = task.turnId;
if (turnId != null) {
if (finishAsyncTaskRaceHook != null) {
@@ -1276,6 +1341,49 @@ public final class MessageService {
this.finishAsyncTaskRaceHook = hook;
}
/**
* Null in production; test seam for fleetd #329 (F2) — invoked from {@link #finishAsyncTask(Task,
* Reply)} unconditionally, immediately after {@code task.future.complete(result)} runs (before
* {@link Task#turnId} is even read, so it fires regardless of whether this task was ever asked).
* A test installs this to force an exception into exactly the shape fleetd #329 identified:
* something throws inside {@link #sendAsync}'s executor task after the async ticket's future is
* already resolved, so the surrounding {@code catch (Throwable t)} can only report failure through
* {@code completeExceptionally} — a silent no-op on an already-completed future. Engineering a
* real exception to land in that exact post-completion window is what the ticket itself had to do
* by temporarily deleting a production guard (fleetd #324's {@code turnId != null} check); this
* hook drives the identical shape deterministically instead.
*/
private volatile Runnable afterFinishAsyncTaskCompleteHookForTest;
/**
* Test-only (fleetd #329, F2): install {@link #afterFinishAsyncTaskCompleteHookForTest}.
* Package-private so the test, in the same package, can reach it without widening any production
* API.
*/
void setAfterFinishAsyncTaskCompleteHookForTest(Runnable hook) {
this.afterFinishAsyncTaskCompleteHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
* value is used as the {@code asyncTasksByTurn} removal key. Mirrors {@link
* #finishAsyncTaskRaceHook} exactly, for the structurally identical double-read fleetd #324 fixed
* in {@link #finishAsyncTask}: a test installs this to force, deterministically, {@code
* clearAsyncQuestion}'s unlocked {@code forgetTurn=true} path nulling {@link Task#turnId} in that
* exact window, and to confirm the single-read fix tolerates it (the captured local is used
* unconditionally, so a hook that nulls the field afterward cannot affect this call).
*/
private volatile Runnable replyOrphanTurnIdRaceHookForTest;
/**
* Test-only (fleetd #329, F3): install {@link #replyOrphanTurnIdRaceHookForTest}. Package-private
* so the test, in the same package, can reach it without widening any production API.
*/
void setReplyOrphanTurnIdRaceHookForTest(Runnable hook) {
this.replyOrphanTurnIdRaceHookForTest = hook;
}
/**
* Test-only (fleetd #324): run the exact production cleanup {@link #ask}'s own timeout path runs
* unlocked — {@link #clearAsyncQuestion(String, boolean)} with {@code forgetTurn=true} — so a test
@@ -1285,14 +1393,6 @@ public final class MessageService {
clearAsyncQuestion(turnId, true);
}
/** Complete the async ticket correlated to a specific answered turn. */
private void finishAsyncTask(String turnId, Reply result) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
finishAsyncTask(task, result);
}
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
@@ -23,6 +23,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.function.Supplier;
@@ -370,6 +371,69 @@ class HerdrPeerLauncherAllowListWiringTest {
"expected the pre-existing 'scrub blanks them' INFO unchanged, got: " + messages);
}
/**
* fleetd #341: {@code unprotectedGapLogged} guarded TWO WARN branches that name DIFFERENT env
* var names — the allow-list branch ({@code keptByDerivedList}, below) and the deny-by-default
* / non-zsh-fallback branch ({@link HerdrPeerLauncher#warnGapUnprotected}). {@code
* memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change between
* two spawns on the same launcher instance — a config reload needs no restart. Spawn 1 runs
* under {@code deny-by-default} with a gap of {@code SPAWN_ONE_UNCOVERED_TOKEN}, which trips
* the (before this fix) SHARED one-shot flag. The policy is then reloaded to {@code
* allow-list}; spawn 2's gap is {@code FLEETD_WORKER_TOKEN} instead — the test profile's own
* {@code tokenEnv}, which the derived allow-list keeps even though it is on neither {@code
* known:} nor {@code allow:}, so it is genuinely unprotected and deserves its own WARN. Before
* this fix that WARN never fires, because the shared flag was already {@code true} — the
* operator is never told {@code FLEETD_WORKER_TOKEN} reaches every member pane unblocked. Real
* path: two real {@link HerdrPeerLauncher#spawn} calls on ONE launcher instance, with mutable
* {@code memberCredentials}/host-env suppliers standing in for a live config reload between
* spawns.
*/
@Test
void aDifferentUnprotectedGapOnALaterSpawnIsNotSuppressedByAnEarlierSpawnsWarn() {
FakeHerdr herdr = new FakeHerdr();
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
new FleetConfig.MemberCredentials(null, List.of(), List.of(), null)); // deny-by-default
AtomicReference<Set<String>> hostEnvState =
new AtomicReference<>(Set.of("SPAWN_ONE_UNCOVERED_TOKEN"));
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.WARN);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
// Spawn 1: deny-by-default, gap = {SPAWN_ONE_UNCOVERED_TOKEN} — the effectiveAllowed ==
// null branch, via warnGapUnprotected.
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
// Live policy reload to allow-list, with a DIFFERENT gap name.
credsState.set(new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null));
hostEnvState.set(Set.of("FLEETD_WORKER_TOKEN"));
// Spawn 2: allow-list, gap = {FLEETD_WORKER_TOKEN} — kept by the derived allow-list
// (the profile's own tokenEnv), so it is the effectiveAllowed != null / keptByDerivedList
// branch, at the SAME log line HerdrPeerLauncher:1819 guards with the shared flag.
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
assertTrue(messages.stream().anyMatch(
m -> m.contains("UNBLOCKED") && m.contains("SPAWN_ONE_UNCOVERED_TOKEN")),
"spawn 1's deny-by-default gap must still warn — got: " + messages);
assertTrue(messages.stream().anyMatch(
m -> m.contains("UNBLOCKED") && m.contains("FLEETD_WORKER_TOKEN")),
"spawn 2's gap names a DIFFERENT env var than spawn 1 (FLEETD_WORKER_TOKEN, not "
+ "SPAWN_ONE_UNCOVERED_TOKEN) — it must still be warned about even though a "
+ "flag already fired once for spawn 1's unrelated name. Before fleetd #341's "
+ "fix this WARN never fires because unprotectedGapLogged was already true. "
+ "Got: " + messages);
}
/**
* fleetd #185 stage 2: with {@code memberHerdrSocket:} configured, member panes run under a
* different OS user — {@link HerdrPeerLauncher#hostEnvNames} describes fleetd's own process, not
@@ -1,5 +1,9 @@
package dev.ltms.fleet.msg;
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.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
@@ -11,6 +15,7 @@ import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.inject.Injector;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@@ -1041,6 +1046,105 @@ class MessageServiceTest {
"the reply completed its own ticket directly and never touched the inbox");
}
/**
* fleetd #329 (F1). {@code answer()} completes the async ticket by looking {@code turnId} up in
* {@code asyncTasksByTurn} a SECOND time (the first is at :991, purely to re-register the
* {@code asyncTasksByWaiter} entry for #282's chained-ask case). That second lookup races
* {@code ask()}'s own unlocked timeout cleanup ({@code clearAsyncQuestion(turnId, true)}, run from
* {@code markAskTimedOut} + forgetting): {@code ask()}'s {@code ticket.answer().get(timeoutMillis)}
* can time out at essentially the same instant {@code answer()}'s {@code rendezvous.answerAsk}
* call above already succeeded and unblocked the worker. fleetd #324's single-read fix does not
* help here — that fixed a torn read of one already-held {@link MessageService} internal
* {@code Task}; this is a second, independent map lookup by a caller that no longer holds the
* {@code Task} it already found once.
*
* <p>This test does not wait for that real race to land on its own schedule — it drives the exact
* sequence the ticket describes (worker asks, primary answers, worker's real reply arrives) and
* fires the identical production cleanup {@code ask()}'s timeout path runs
* ({@code clearAsyncQuestion(turnId, true)}, via the {@code forgetTurnForTest} seam fleetd #324
* already merged) at the point between "the worker's own {@code ask()} call has unblocked" and
* "the worker's real {@code fleet_reply} arrives" — the exact window fleetd #329 names.
*
* <p>What this proves: given that exact interleaving, the async ticket must still resolve to the
* worker's real reply, not stay stuck {@link MessageService.Phase#PENDING} forever. What it does
* not prove: that the interleaving itself is reachable in production on its own timing — that is
* established by reading the code (see the ticket's "path in"), not by this test, for the same
* reason fleetd #324's own race test says so.
*/
@Test
void aReplyRacingAsksTimeoutCleanupStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer(),
"the worker's own ask() call must have already unblocked with the primary's answer "
+ "before we force the race below");
// ask()'s own unlocked timeout cleanup can forget this exact turnId at essentially the same
// instant answer() has already unblocked the worker and is now waiting on the resumed turn's
// 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"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the primary's own answer() call must still see the worker's real reply");
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"fleet_poll{ticket} must return the worker's real reply, not stay PENDING forever "
+ "just because ask()'s timeout cleanup forgot this turnId first");
}
/**
* fleetd #329 (F3). {@link MessageService#reply} reads the same orphan {@code Task}'s {@code
* turnId} twice on this path (once to check it is non-null, once as the {@code
* asyncTasksByTurn.remove} key) — the exact double-read shape fleetd #324 fixed in {@code
* finishAsyncTask}. This test drives the same real sequence as {@code
* aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket} (worker asks, primary answers, the
* primary's own bounded wait for the resumed turn expires) to reach a task with {@code turnId}
* genuinely stamped and no live rendezvous waiter open — the state {@code askAnsweredAsyncTasks}
* matches here — then fires the identical production cleanup {@code ask()}'s own timeout path
* runs ({@code clearAsyncQuestion(turnId, true)}, via the {@code forgetTurnForTest} seam fleetd
* #324 merged) at the point between the check and the removal use, via a dedicated test hook
* mirroring {@code finishAsyncTaskRaceHook}.
*/
@Test
void replyToAnOrphanedTaskSurvivesTurnIdGoingNullBetweenItsTwoReads() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait must give up first, leaving turnId stamped with no "
+ "live waiter — the state reply()'s F3 code path matches");
messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
try {
assertTrue(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(),
"the ticket must still resolve to the worker's real reply despite the forced race");
} finally {
messages.setReplyOrphanTurnIdRaceHookForTest(null);
}
}
/**
* fleetd #307's ambiguity guard: an ask timeout frees its target ({@code hasAsyncQuestion}
* becomes false the instant it lapses — proven above), so a second, independent delegation can
@@ -1151,6 +1255,63 @@ class MessageServiceTest {
assertEquals("done", done.reply());
}
// --- fleetd #329 (F2): an exception after the ticket's future completes must reach a log -----
//
// sendAsync's own executor task ends with catch (Throwable t) { task.future.completeExceptionally(t); }
// — but finishAsyncTask completes that same future on its first line, so anything that throws
// afterward hits an already-completed future: completeExceptionally returns false and does
// nothing, and (before this fix) nothing logged it either. Measured on the ticket: temporarily
// reintroducing the fleetd #324 NPE reproduced 19 real exceptions on the ordinary path with
// 74/74 tests staying green and zero log lines. There is no reachable production call site where
// finishAsyncTask throws after completing the future (fleetd #324 already closed the one that
// used to), so this test drives the shape directly via a dedicated test-only hook
// (afterFinishAsyncTaskCompleteHookForTest) rather than trying to engineer a real exception into
// that narrow window — the same technique fleetd #324's own race test and this ticket's F1/F3
// tests use for their own races.
private static ListAppender<ILoggingEvent> attachMessageServiceLog() {
Logger logger = (Logger) LoggerFactory.getLogger(MessageService.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detachMessageServiceLog(ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(MessageService.class)).detachAppender(appender);
}
@Test
void anExceptionAfterTheTicketFutureCompletesStillReachesTheLog() throws Exception {
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try {
messages.setAfterFinishAsyncTaskCompleteHookForTest(() -> {
throw new RuntimeException("PROBE-329-F2");
});
String ticket = messages.sendAsync(T, "do the task");
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
// finishAsyncTask's own first line already completed the future with the real result
// BEFORE the hook threw — F2 is a visibility gap, not a correctness gap for the ticket
// itself, so the ticket's own outcome must be unaffected by the swallowed exception.
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("done", done.reply());
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.ERROR
&& e.getFormattedMessage().contains(ticket)
&& e.getThrowableProxy() != null
&& "PROBE-329-F2".equals(e.getThrowableProxy().getMessage())),
"an exception thrown after the ticket's future already completed must still "
+ "reach the log, not vanish silently");
} finally {
messages.setAfterFinishAsyncTaskCompleteHookForTest(null);
detachMessageServiceLog(appender);
}
}
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
@Test