Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 667254df47 | |||
| 77ad88631b | |||
| d88017807b | |||
| 2159a5a94a |
@@ -385,9 +385,12 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
}
|
||||
}
|
||||
|
||||
// unreachable.size() counts DISTINCT profiles, not attempts (a HashSet dedupes a profile
|
||||
// added twice) — say "distinct" so the count matches the sentence and the profile list that
|
||||
// follows, rather than reading as a count of attempts made (fleetd #315).
|
||||
throw new PeerUnreachableException(
|
||||
"no reachable worker profile available after trying " + unreachable.size()
|
||||
+ " candidate(s): " + String.join(", ", unreachable));
|
||||
+ " distinct candidate(s): " + String.join(", ", unreachable));
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -5,10 +5,14 @@ import java.util.List;
|
||||
|
||||
/**
|
||||
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
|
||||
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps and
|
||||
* reachability so that a pre-existing config behaves identically after upgrade.
|
||||
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps
|
||||
* ({@code maxLoad}) so that a pre-existing config behaves identically after upgrade — capacity
|
||||
* gating for automatic placement is deliberately out of scope for {@code fixed}, exactly as it
|
||||
* always has been. Reachability is a narrower exception (fleetd #315, below): a profile is never
|
||||
* checked for reachability up front, only skipped once it has already failed in <em>this same</em>
|
||||
* spawn call's retry loop — see the unreachable case below.
|
||||
*
|
||||
* <p>Three exceptions walk past the default instead of returning it unconditionally:
|
||||
* <p>Four exceptions walk past the default instead of returning it unconditionally:
|
||||
* <ul>
|
||||
* <li>Quarantine (CB-578 stage B): a quarantined default is a credential that just refused on
|
||||
* a usage limit, not a transient capacity or reachability concern.
|
||||
@@ -16,13 +20,21 @@ import java.util.List;
|
||||
* ({@code BackendOutagePolicy}) — a separate, shorter-lived source from quarantine. When a
|
||||
* profile is both quarantined and cooling off, only the quarantine reason is reported
|
||||
* (exhaustion takes priority), matching {@code CompositePeerLauncher}'s explicit-spawn order.
|
||||
* <li>Unreachable (fleetd #315): {@code CompositePeerLauncher.spawn} retries a failed candidate
|
||||
* on the next one and rebuilds the {@link PlacementContext} so {@code ctx.unreachable()}
|
||||
* names every profile that already failed with {@code PeerUnreachableException} in this same
|
||||
* call. Without this check {@code select} kept handing back the same dead default forever —
|
||||
* the retry loop's own comment says "so the policy excludes this profile", and this is what
|
||||
* makes that true for {@code fixed} too, matching {@code weighted}/{@code round-robin}
|
||||
* (both filter on {@code ctx.unreachable()} via {@link PlacementPolicyUtil#available}).
|
||||
* <li>Weight 0 (CB-554): {@code fixed} is still automatic selection, so a profile the operator
|
||||
* marked "never auto-select me" ({@code weight <= 0}) must be skipped here exactly as
|
||||
* {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming
|
||||
* the profile is unaffected, only this automatic fallback walk.
|
||||
* </ul>
|
||||
* A fleet where nothing is ever quarantined, cooling off, or weight-0 never exercises any of these
|
||||
* paths, so today's behaviour is unchanged.
|
||||
* A fleet where nothing is ever quarantined, cooling off, unreachable, or weight-0 never exercises
|
||||
* any of these paths, so today's behaviour is unchanged — in particular, the very first selection
|
||||
* of a spawn call always sees an empty {@code unreachable} set, so the first choice is untouched.
|
||||
*/
|
||||
final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
|
||||
@@ -30,12 +42,12 @@ final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
public PlacementCandidate select(PlacementContext ctx) {
|
||||
String d = ctx.defaultProfile();
|
||||
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d) && !ctx.coolingOff().contains(d)
|
||||
&& !weightExcluded(ctx, d)) {
|
||||
&& !ctx.unreachable().contains(d) && !weightExcluded(ctx, d)) {
|
||||
return new PlacementCandidate(d, null, 1.0f, null);
|
||||
}
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
if (!ctx.quarantined().contains(c.profile()) && !ctx.coolingOff().contains(c.profile())
|
||||
&& !c.excluded()) {
|
||||
&& !ctx.unreachable().contains(c.profile()) && !c.excluded()) {
|
||||
return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad());
|
||||
}
|
||||
}
|
||||
@@ -44,8 +56,9 @@ final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
// Exhaustion quarantine takes priority: reported only when quarantine is absent, so the
|
||||
// message never claims "cooling off" for a profile that is really backend-exhausted.
|
||||
boolean dCoolingOff = !dQuarantined && ctx.coolingOff().contains(d);
|
||||
boolean dUnreachable = ctx.unreachable().contains(d);
|
||||
boolean dWeightExcluded = weightExcluded(ctx, d);
|
||||
if (dQuarantined || dCoolingOff || dWeightExcluded) {
|
||||
if (dQuarantined || dCoolingOff || dUnreachable || dWeightExcluded) {
|
||||
List<String> reasons = new ArrayList<>();
|
||||
if (dQuarantined) {
|
||||
reasons.add("is quarantined (backend exhausted)");
|
||||
@@ -53,6 +66,9 @@ final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
if (dCoolingOff) {
|
||||
reasons.add("is cooling off after repeated backend errors");
|
||||
}
|
||||
if (dUnreachable) {
|
||||
reasons.add("is unreachable");
|
||||
}
|
||||
if (dWeightExcluded) {
|
||||
reasons.add("has weight 0 (excluded from automatic selection)");
|
||||
}
|
||||
@@ -62,7 +78,7 @@ final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
}
|
||||
if (!ctx.candidates().isEmpty()) {
|
||||
throw new PlacementException("all worker profiles are excluded from automatic "
|
||||
+ "selection (quarantined, cooling off, or weight-0)");
|
||||
+ "selection (quarantined, cooling off, unreachable, or weight-0)");
|
||||
}
|
||||
throw new PlacementException("no worker profiles configured");
|
||||
}
|
||||
|
||||
@@ -386,6 +386,31 @@ public final class SessionManager implements TurnListener {
|
||||
// from the registry with no pane stop is an orphaned pane — a live terminal burning a fleet
|
||||
// slot that no longer appears in the roster and can never be reclaimed.
|
||||
launcher.stop(paneId);
|
||||
if (removed != null && !preserveWorktree && removed.worktree() != null) {
|
||||
// fleetd #316: the `dirty` read above ran while the worker could still write to this
|
||||
// worktree, so a stale `false` must not be trusted to authorise the --force removal
|
||||
// below. Re-read the worktree's state one more time, right here — immediately before
|
||||
// the one step that would destroy it, and only on the path that is actually about to
|
||||
// do that (invariant 4: no second unconditional `git status` on a release that already
|
||||
// decided to preserve). By now `launcher.stop` has returned, so this read reflects
|
||||
// whatever the worker managed to write up to and including its teardown, not whatever
|
||||
// it had written at release-start time.
|
||||
if (dirtyImmediatelyBeforeRemoval(removed)) {
|
||||
preserveWorktree = true;
|
||||
// The pre-stop snapshot above never ran for this session (the pre-stop read said
|
||||
// clean), so this is the only chance to get the newly-discovered work into
|
||||
// refs/wip/* rather than leaving the on-disk preserve as the sole copy. Best-effort,
|
||||
// like every other snapshot attempt — trySnapshot logs and swallows its own failure.
|
||||
String lateSnapshotRef = trySnapshot(removed, cause);
|
||||
log.warn("release {} preserves worktree {} for pane={} terminal={}: it reported "
|
||||
+ "clean before the pane stopped but dirty immediately before removal — the "
|
||||
+ "worker wrote to it during teardown, and --force removing it now would "
|
||||
+ "have destroyed that work{}",
|
||||
cause, removed.worktree(), paneId, removed.terminalId(),
|
||||
lateSnapshotRef == null ? "" : " (snapshotted to refs/wip/" + removed.branch()
|
||||
+ " commit=" + lateSnapshotRef + ")");
|
||||
}
|
||||
}
|
||||
if (removed != null && !preserveWorktree && removed.worktree() != null) {
|
||||
// fleetd #283: this is the one cleanup step in this method that used to be bare. By the
|
||||
// time it runs, the registry entry, the retained handle, and the pane are all already
|
||||
@@ -404,6 +429,24 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #316: the read that actually authorises {@code worktrees.remove}, taken with the
|
||||
* worker's pane already stopped. Fails toward preserving (returns {@code true}) on any
|
||||
* exception — the same rule the pre-stop check applies (CB-581): once we can no longer tell
|
||||
* whether the worktree is dirty, preserving costs disk while deleting on a guess can destroy
|
||||
* work that has no other copy.
|
||||
*/
|
||||
private boolean dirtyImmediatelyBeforeRemoval(MemberSession removed) {
|
||||
try {
|
||||
return worktrees.hasUncommitted(removed.worktree());
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("release could not re-check worktree {} for pane={} terminal={} immediately "
|
||||
+ "before removal; preserving it rather than risk destroying unsaved work: {}",
|
||||
removed.worktree(), removed.paneId(), removed.terminalId(), e.toString());
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Best-effort snapshot of a dirty worktree into {@code refs/wip/<branch>} (CB-578 stage C). A
|
||||
* failure here must never escalate: the caller has already decided to preserve the worktree
|
||||
|
||||
@@ -615,6 +615,67 @@ class CompositePeerLauncherTest {
|
||||
assertEquals(1, adapter.spawnCount("b"));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #315: {@code CompositePeerLauncher.spawn} rebuilds the {@link PlacementContext} after
|
||||
* every failed attempt "so the policy excludes this profile" (see the comment at the retry call
|
||||
* site) — but {@code FixedPlacementPolicy} never read {@code ctx.unreachable()}, so under the
|
||||
* default {@code fixed} placement every retry re-picked the same dead default and a second,
|
||||
* healthy, configured profile was never tried. This is the same scenario as
|
||||
* {@link #failoverRetriesNextCandidateWhenProfileIsUnreachable}, but pinned to {@code fixed()}
|
||||
* instead of {@code weighted()} — the three existing failover tests all use {@code weighted()},
|
||||
* which is exactly why nobody caught this: the retry loop's contract has no coverage under its
|
||||
* own default policy.
|
||||
*/
|
||||
@Test
|
||||
void failoverRetriesNextCandidateUnderFixedPlacementWhenProfileIsUnreachable() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"a", stubWorker("a"),
|
||||
"b", stubWorker("b"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of("a"));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "a", profiles, PlacementPolicies.fixed(), _ -> 0);
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertEquals("b", h.profile(),
|
||||
"fixed placement must fail over from the unreachable default a to the healthy b");
|
||||
assertEquals(1, adapter.spawnCount("a"), "a was tried once and failed");
|
||||
assertEquals(1, adapter.spawnCount("b"), "b was tried once and succeeded");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #315: the same fix — {@code FixedPlacementPolicy} consulting {@code ctx.unreachable()}
|
||||
* — also covers the wiring-bug branch in {@code CompositePeerLauncher.spawn}: a profile that
|
||||
* placement is allowed to choose (it is in the configured candidate list) but that no delegate
|
||||
* declares ({@code byProfile.get(chosen.profile()) == null}). That branch adds the profile to
|
||||
* {@code unreachable} and {@code continue}s without ever calling a launcher, so before this fix
|
||||
* {@code fixed} handed back the same adapterless profile on every remaining attempt too.
|
||||
*/
|
||||
@Test
|
||||
void failoverSkipsAConfiguredProfileNoAdapterDeclaresUnderFixedPlacement() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
// Placement's candidate list has three profiles, in this order (LinkedHashMap preserves it,
|
||||
// and the fixed default resolves to the first — see the `ordered` helper's own javadoc).
|
||||
Map<String, FleetConfig.Profile> profiles = new LinkedHashMap<>();
|
||||
profiles.put("c", stubWorker("c"));
|
||||
profiles.put("a", stubWorker("a"));
|
||||
profiles.put("b", stubWorker("b"));
|
||||
// The adapter only declares a and b — c is a configured profile with no owning adapter,
|
||||
// the "wiring bug" the comment in CompositePeerLauncher.spawn calls out.
|
||||
Map<String, FleetConfig.Profile> adapterProfiles = new LinkedHashMap<>();
|
||||
adapterProfiles.put("a", stubWorker("a"));
|
||||
adapterProfiles.put("b", stubWorker("b"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, adapterProfiles, "a", Set.of());
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "a", profiles, PlacementPolicies.fixed(), _ -> 0);
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertEquals("a", h.profile(),
|
||||
"c has no adapter, so fixed placement must skip it and land on the next candidate, a");
|
||||
assertEquals(0, adapter.spawnCount("c"), "c is never spawned — no adapter owns it");
|
||||
assertEquals(1, adapter.spawnCount("a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void explicitSpawnAtMaxLoadThrowsPlacementExceptionNamingProfileLiveAndCap() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -86,12 +86,32 @@ class SessionManagerTest {
|
||||
private volatile RuntimeException hasUncommittedFailure;
|
||||
private volatile RuntimeException snapshotFailure;
|
||||
private final java.util.concurrent.atomic.AtomicLong snapshotSeq = new java.util.concurrent.atomic.AtomicLong();
|
||||
/** fleetd #316: successive {@code hasUncommitted} answers, one per call, last one sticky
|
||||
* once exhausted — models a worktree whose state changes between reads. Empty (the
|
||||
* default) falls back to the plain {@link #dirty} flag, so every existing test using this
|
||||
* fake keeps returning one fixed answer. */
|
||||
private final List<Boolean> dirtySequence = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
private final java.util.concurrent.atomic.AtomicInteger hasUncommittedCalls =
|
||||
new java.util.concurrent.atomic.AtomicInteger();
|
||||
|
||||
RecordingWorktrees dirty(boolean dirty) {
|
||||
this.dirty = dirty;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** fleetd #316: return {@code answers[0]} on the first {@code hasUncommitted} call,
|
||||
* {@code answers[1]} on the second, and so on; the last element repeats after that. */
|
||||
RecordingWorktrees dirtySequence(boolean... answers) {
|
||||
for (boolean a : answers) {
|
||||
dirtySequence.add(a);
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
int hasUncommittedCallCount() {
|
||||
return hasUncommittedCalls.get();
|
||||
}
|
||||
|
||||
RecordingWorktrees failHasUncommittedWith(RuntimeException e) {
|
||||
this.hasUncommittedFailure = e;
|
||||
return this;
|
||||
@@ -127,9 +147,13 @@ class SessionManagerTest {
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
int call = hasUncommittedCalls.getAndIncrement();
|
||||
if (hasUncommittedFailure != null) {
|
||||
throw hasUncommittedFailure;
|
||||
}
|
||||
if (!dirtySequence.isEmpty()) {
|
||||
return dirtySequence.get(Math.min(call, dirtySequence.size() - 1));
|
||||
}
|
||||
return dirty;
|
||||
}
|
||||
|
||||
@@ -1207,6 +1231,80 @@ class SessionManagerTest {
|
||||
+ "dirty check threw");
|
||||
}
|
||||
|
||||
// --- fleetd #316: the dirty check must be re-taken after the worker is stopped, not trusted
|
||||
// stale from before it ------------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void releaseDoesNotRemoveAWorktreeThatBecameDirtyBetweenTheFirstCheckAndRemoval() {
|
||||
// Models the exact race #316 reports: hasUncommitted answers clean while the worker is
|
||||
// still running (call 1), the worker then writes new work, and by the time release is
|
||||
// about to force-remove the worktree a second read (call 2) would see it as dirty. Without
|
||||
// the fix this test fails: release() never re-reads and force-removes the worktree anyway.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirtySequence(false, true);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316a", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"a worktree that turned dirty between the pre-stop read and removal must be preserved");
|
||||
assertEquals(2, worktrees.hasUncommittedCallCount(),
|
||||
"the fix re-reads hasUncommitted exactly once more, immediately before removal");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseSnapshotsWorkFoundOnlyByTheLateRecheck() {
|
||||
// #316's second half: the pre-stop dirty=false means trySnapshot never ran for this
|
||||
// session, so the late-discovered work would otherwise have no refs/wip/* copy at all —
|
||||
// only the on-disk preserve. The re-check path must snapshot it too.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirtySequence(false, true);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316b", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(java.util.List.of(s.worktree()), worktrees.snapshotCalls(),
|
||||
"the newly-dirty worktree is snapshotted even though the pre-stop check saw it clean");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseStillRemovesAWorktreeThatStaysCleanOnTheLateRecheck() {
|
||||
// The ordinary, non-racing case: nothing else changes behaviour when the second read
|
||||
// agrees with the first.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(false);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316c", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(java.util.List.of(s.worktree()), worktrees.removeCalls(),
|
||||
"a worktree that is still clean on the late recheck is removed as before");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseNeverReChecksAWorktreeAlreadyPreservedByTheFirstDirtyCheck() {
|
||||
// Invariant 4 from #316: no second unconditional git status. A release that already
|
||||
// decided to preserve (the ordinary CB-576 dirty path) must not pay for a second read.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(true);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316d", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(1, worktrees.hasUncommittedCallCount(),
|
||||
"a release that already preserves on the first read must not re-check before "
|
||||
+ "skipping the removal it was never going to do");
|
||||
assertTrue(worktrees.removeCalls().isEmpty());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #283 defect 1 changed this test's own premise, so its assertions are updated along
|
||||
* with the production fix. Before #283, the middle session's worktree-removal failure escaped
|
||||
|
||||
Reference in New Issue
Block a user