Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d88017807b | |||
| 2159a5a94a | |||
| b9c2cf69f4 | |||
| 8beae50fe7 | |||
| b2a58cb966 | |||
| f159ca7d27 | |||
| 83f2aea60f | |||
| 002329adb5 | |||
| a49671ceb9 |
@@ -22,6 +22,7 @@ import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.ShuttingDownException;
|
||||
import dev.ltms.fleet.session.WorktreeRequest;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
@@ -962,6 +963,10 @@ public final class FleetMcp {
|
||||
return text(json(memberView(member)));
|
||||
} catch (GuardException e) {
|
||||
return error("subscription boundary: " + e.getMessage());
|
||||
} catch (ShuttingDownException e) {
|
||||
// fleetd #308: the daemon's shutdown drain has already started — refuse loudly rather
|
||||
// than register a session drainAll will never see again.
|
||||
return error("shutting down: " + e.getMessage());
|
||||
} catch (PlacementException e) {
|
||||
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct
|
||||
// from "profile does not exist" below.
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
|
||||
@@ -19,6 +19,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.ShuttingDownException;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.WorktreeRequest;
|
||||
@@ -501,6 +502,11 @@ public final class FleetApp {
|
||||
ctx.status(201).json(view(member));
|
||||
} catch (GuardException e) {
|
||||
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
|
||||
} catch (ShuttingDownException e) {
|
||||
// fleetd #308: the daemon's shutdown drain has already started — 503, not a bare 500,
|
||||
// so this reads the same as PlacementException below: valid request, refused because
|
||||
// of a transient daemon state rather than a bad argument.
|
||||
ctx.status(503).json(Map.of("error", "shutting_down", "detail", e.getMessage()));
|
||||
} catch (PlacementException e) {
|
||||
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign,
|
||||
// likely-transient refusal, distinct from "profile does not exist" below. 503: the
|
||||
|
||||
@@ -93,6 +93,9 @@ public final class GitWorktrees implements Worktrees {
|
||||
/** OS group name for {@link #shareWithGroup} (fleetd #185 stage 3); {@code null} ⇒ feature off. */
|
||||
private final String group;
|
||||
private final Consumer<String> afterWorktreeAdded;
|
||||
/** How the initial {@code git worktree add} command runs. Package-private test seam for an
|
||||
* interrupted command after Git has made worktree state. */
|
||||
private final Function<String[], String> worktreeAddRunner;
|
||||
/** How {@link #shareWithGroup}'s processes (git config / chgrp / chmod / find) actually run.
|
||||
* Defaults to the real {@link #exec(String...)}. Package-private test seam so a unit test can
|
||||
* prove "no group configured ⇒ zero processes spawned" and inspect exactly what a configured
|
||||
@@ -132,7 +135,13 @@ public final class GitWorktrees implements Worktrees {
|
||||
|
||||
/** Test seam combining a configurable {@code group} with {@link #afterWorktreeAdded}. */
|
||||
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded) {
|
||||
this(configuredRoot, group, afterWorktreeAdded, null);
|
||||
this(configuredRoot, group, afterWorktreeAdded, null, null);
|
||||
}
|
||||
|
||||
/** Test seam for changing how {@link #shareWithGroup}'s processes run. */
|
||||
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
|
||||
Function<String[], String> shareGroupRunner) {
|
||||
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -141,13 +150,15 @@ public final class GitWorktrees implements Worktrees {
|
||||
* exactly what commands a configured group runs, without a real second OS user/group.
|
||||
*
|
||||
* @param shareGroupRunner {@code null} ⇒ the real {@link #exec(String...)}.
|
||||
* @param worktreeAddRunner {@code null} ⇒ the real {@link #exec(String...)}.
|
||||
*/
|
||||
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
|
||||
Function<String[], String> shareGroupRunner) {
|
||||
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner) {
|
||||
this.configuredRoot = configuredRoot;
|
||||
this.group = (group == null || group.isBlank()) ? null : group;
|
||||
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
|
||||
this.shareGroupRunner = shareGroupRunner != null ? shareGroupRunner : this::exec;
|
||||
this.worktreeAddRunner = worktreeAddRunner != null ? worktreeAddRunner : this::exec;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -166,8 +177,8 @@ public final class GitWorktrees implements Worktrees {
|
||||
String wt = path.toAbsolutePath().toString();
|
||||
log.info("adding worktree branch={} path={} base={}", branch, wt, base);
|
||||
removeUserInfoFromHttpsOrigin(repoRoot);
|
||||
exec("git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base);
|
||||
try {
|
||||
worktreeAddRunner.apply(new String[] {"git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base});
|
||||
afterWorktreeAdded.accept(wt);
|
||||
requireCredentialFreeHttpsOrigin(wt);
|
||||
configureEnvironmentCredentialHelper(repoRoot, wt);
|
||||
@@ -181,10 +192,11 @@ public final class GitWorktrees implements Worktrees {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code add()} has already created the worktree and its branch by the time any step from
|
||||
* {@link #afterWorktreeAdded} through {@link #isolateToolSurface} can throw — including
|
||||
* {@code git worktree add} may have created the worktree and its branch by the time it, or any
|
||||
* later step through {@link #isolateToolSurface}, throws. This includes
|
||||
* {@link #requireCredentialFreeHttpsOrigin}, an intended security refusal, not only an IO
|
||||
* accident. Without this, {@code add()} never returns, so its caller
|
||||
* accident. A Git-reported {@code worktree add} failure usually creates nothing, but an
|
||||
* interrupted command can leave partial state. Without cleanup, {@code add()} never returns, so its caller
|
||||
* ({@code SessionManager#acquireWithWorktree}) never receives a path to register or clean up:
|
||||
* its local {@code path} stays null, the {@code if (path != null)} guard in its own catch block
|
||||
* never runs, and the worktree directory and branch leak on disk forever with nothing tracking
|
||||
@@ -204,6 +216,10 @@ public final class GitWorktrees implements Worktrees {
|
||||
* used in {@code SessionManager#acquireWithWorktree}'s own catch block.
|
||||
*/
|
||||
private void cleanupAfterAddFailure(String repoRoot, String worktreePath, String branch, RuntimeException original) {
|
||||
if (!Files.exists(Path.of(worktreePath))) {
|
||||
log.debug("provisioning failed before worktree {} existed; nothing to clean up", worktreePath);
|
||||
return;
|
||||
}
|
||||
log.warn("provisioning failed for branch={} path={}: {} — cleaning up before rethrowing",
|
||||
branch, worktreePath, original.getMessage());
|
||||
try {
|
||||
|
||||
@@ -21,6 +21,7 @@ import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -61,6 +62,8 @@ public final class SessionManager implements TurnListener {
|
||||
private final LongSupplier nowNanos;
|
||||
private final int contextCap;
|
||||
private final boolean clearAfterTurn;
|
||||
/** Null in production; test seam for the interval before an idle session's conditional release. */
|
||||
private final Consumer<MemberSession> beforeIdleRelease;
|
||||
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
|
||||
/**
|
||||
* CB-586: the repo root the fleet actually works in, remembered the first time a worktree
|
||||
@@ -71,6 +74,16 @@ public final class SessionManager implements TurnListener {
|
||||
*/
|
||||
private volatile String fleetRepoRoot;
|
||||
|
||||
/**
|
||||
* fleetd #308: flips true the instant {@link #drainAll} starts, before its registry snapshot
|
||||
* is even taken — so a spawn already in flight sees the refusal as early as a plain flag can
|
||||
* make it. This alone cannot close the race completely: a caller that read {@code false} just
|
||||
* before the flip can still land in the registry after the snapshot. {@link #drainAll}'s
|
||||
* post-loop sweep is what catches that straggler; the two mechanisms are deliberately paired,
|
||||
* see {@link #drainAll}'s javadoc.
|
||||
*/
|
||||
private final AtomicBoolean draining = new AtomicBoolean(false);
|
||||
|
||||
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
|
||||
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
/** CB-516: notified with a {@link ReleaseDetail} on every release; no-op until wired. */
|
||||
@@ -102,13 +115,23 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
|
||||
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
|
||||
int contextCap, boolean clearAfterTurn) {
|
||||
int contextCap, boolean clearAfterTurn) {
|
||||
this(launcher, worktrees, nowNanos, contextCap, clearAfterTurn, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Package-private constructor for a deterministic reap/delivery race test. Production callers
|
||||
* use the constructor above, whose null hook adds no callback or lock to an ordinary reap.
|
||||
*/
|
||||
SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
|
||||
int contextCap, boolean clearAfterTurn, Consumer<MemberSession> beforeIdleRelease) {
|
||||
this.launcher = launcher;
|
||||
this.worktrees = worktrees;
|
||||
this.presence = new PresenceFleet(this);
|
||||
this.nowNanos = nowNanos;
|
||||
this.contextCap = contextCap;
|
||||
this.clearAfterTurn = clearAfterTurn;
|
||||
this.beforeIdleRelease = beforeIdleRelease;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -181,6 +204,14 @@ public final class SessionManager implements TurnListener {
|
||||
public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd,
|
||||
String ownerTerminal, WorktreeRequest wt,
|
||||
String sessionName, String resumeSessionId) {
|
||||
// fleetd #308: refuse before anything else runs — no slot reservation, no launcher spawn —
|
||||
// so a caller learns the daemon is going down instead of getting a session drainAll will
|
||||
// never see again. Checked here because every other acquire(...) overload delegates to
|
||||
// this one, so this is the single point every spawn path passes through.
|
||||
if (draining.get()) {
|
||||
throw new ShuttingDownException("fleetd is shutting down; refusing to spawn a session "
|
||||
+ "the shutdown drain would never see");
|
||||
}
|
||||
MemberRole memberRole = (role == null) ? MemberRole.DEV : role;
|
||||
requireResumeCapability(profile, resumeSessionId);
|
||||
// CB-619 / fleetd #123: an explicit profile bypasses placement (CompositePeerLauncher only
|
||||
@@ -272,10 +303,32 @@ public final class SessionManager implements TurnListener {
|
||||
*/
|
||||
private void release(String paneId, ReleaseCause cause) {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
|
||||
}
|
||||
|
||||
/**
|
||||
* Tear a session down only while {@code expected} is still its registry value. A lifecycle
|
||||
* transition replaces the immutable record, so this prevents a reap based on an old READY or
|
||||
* DONE record from stopping a worker that delivery has made BUSY.
|
||||
*/
|
||||
private boolean releaseIfCurrent(MemberSession expected, ReleaseCause cause) {
|
||||
if (!registry.remove(expected.paneId(), expected)) {
|
||||
// A lifecycle transition replaced the record between the caller's check and this remove.
|
||||
// Log it: this race is by definition unobservable otherwise, and a reaper that silently
|
||||
// declines to reap is the hardest kind of behaviour to diagnose after the fact.
|
||||
log.debug("skipping reap of pane={}: its registry record changed after the idle check "
|
||||
+ "(most likely a delivery made it BUSY)", expected.paneId());
|
||||
return false;
|
||||
}
|
||||
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
|
||||
return true;
|
||||
}
|
||||
|
||||
private void releaseRemoved(String paneId, MemberSession removed, PeerHandle removedHandle,
|
||||
ReleaseCause cause) {
|
||||
// fleetd #209: remove right alongside the registry entry so a released session's handle is
|
||||
// never leaked — but keep the local reference below, so the id can still be resolved for
|
||||
// the ReleaseDetail this teardown notifies with.
|
||||
PeerHandle removedHandle = handles.remove(paneId);
|
||||
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
|
||||
String snapshotRef = null;
|
||||
if (removed != null) {
|
||||
@@ -846,14 +899,18 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
long idleNanos = now - s.lastActivityAtNanos();
|
||||
if (idleNanos > idleTtlNanos) {
|
||||
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
|
||||
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
|
||||
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
|
||||
// CB-581: one session that fails to release must not abort the whole reaping pass —
|
||||
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
|
||||
try {
|
||||
release(s.paneId());
|
||||
reaped++;
|
||||
if (beforeIdleRelease != null) {
|
||||
beforeIdleRelease.accept(s);
|
||||
}
|
||||
if (releaseIfCurrent(s, ReleaseCause.COMPLETED)) {
|
||||
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
|
||||
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
|
||||
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
|
||||
reaped++;
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
|
||||
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
|
||||
@@ -884,10 +941,40 @@ public final class SessionManager implements TurnListener {
|
||||
* a reason to delete a worker's only copy of its uncommitted work. A session still {@code BUSY}
|
||||
* when the timeout expired is abandoned mid-turn and logged loudly so an operator can find its
|
||||
* kept worktree.
|
||||
*
|
||||
* <p>fleetd #308: {@code roster()} is a one-shot snapshot (see its javadoc), and nothing used
|
||||
* to stop a new session from registering after it was taken — {@link #acquire} stayed open for
|
||||
* as long as this drain waited on a {@code BUSY} session, up to the whole {@code timeoutNanos}
|
||||
* budget. Two things close that window, deliberately paired because neither alone is complete:
|
||||
* {@link #draining} is flipped true before the snapshot is even taken, so {@link #acquire}
|
||||
* refuses (invariant 3: loudly, via {@link ShuttingDownException}) as much of the window as a
|
||||
* plain flag can close; and the sweep below re-reads the registry once the initial snapshot has
|
||||
* fully drained and drains whatever a straggler — a caller that read the flag as {@code false}
|
||||
* a moment before it flipped — still managed to register. The sweep shares the same
|
||||
* {@code deadline} rather than getting its own: {@code timeoutNanos} is a budget for the WHOLE
|
||||
* drain (see above), and a straggler must not buy the drain more time than the flag it lost the
|
||||
* race against would have. In the ordinary case the sweep finds nothing and costs one empty
|
||||
* {@link #roster()} call.
|
||||
*/
|
||||
void drainAll(long timeoutNanos) {
|
||||
long deadline = System.nanoTime() + timeoutNanos;
|
||||
for (MemberSession s : roster()) {
|
||||
draining.set(true);
|
||||
drainSnapshot(roster(), deadline);
|
||||
List<MemberSession> stragglers = roster();
|
||||
if (!stragglers.isEmpty()) {
|
||||
log.warn("drain sweep found {} session(s) registered after the drain snapshot was "
|
||||
+ "taken (raced past the shutdown guard); draining them too", stragglers.size());
|
||||
drainSnapshot(stragglers, deadline);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Drain exactly the sessions in {@code snapshot}, waiting out a {@code BUSY} one against the
|
||||
* shared whole-drain {@code deadline} before releasing it. Shared by {@link #drainAll}'s main
|
||||
* pass and its post-loop straggler sweep (fleetd #308) so both honor the same one budget.
|
||||
*/
|
||||
private void drainSnapshot(List<MemberSession> snapshot, long deadline) {
|
||||
for (MemberSession s : snapshot) {
|
||||
try {
|
||||
if (s.state() == MemberSession.State.BUSY) {
|
||||
while (System.nanoTime() < deadline) {
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
/**
|
||||
* Thrown by {@link SessionManager#acquire} when a spawn is requested after the daemon's shutdown
|
||||
* drain has already begun (fleetd #308).
|
||||
*
|
||||
* <p>{@link SessionManager#drainAll} snapshots the registry once and tears down exactly what is
|
||||
* in that snapshot. A session registered after the snapshot is invisible to the drain loop: its
|
||||
* pane is left running and its worktree is never preserved, and nothing else ever reclaims
|
||||
* either — the daemon's in-memory registry dies with the process. Refusing the spawn here,
|
||||
* loudly, is what stops that session from ever being created in the first place, rather than
|
||||
* silently handing the caller a session the daemon can no longer manage.
|
||||
*/
|
||||
public final class ShuttingDownException extends RuntimeException {
|
||||
public ShuttingDownException(String message) {
|
||||
super(message);
|
||||
}
|
||||
}
|
||||
@@ -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();
|
||||
|
||||
@@ -403,6 +403,52 @@ class GitWorktreesTest {
|
||||
assertTrue(heads.isBlank(), "the branch leaked after a post-creation step threw:\n" + heads);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #309. The add runner is a narrow seam for the case where Git has created state but
|
||||
* the caller then kills the process. The runner first performs the real add in this throwaway
|
||||
* repo, then throws the same kind of exception that {@link GitWorktrees#exec} uses for a timeout.
|
||||
* This proves the failure path cleans both real Git objects without waiting for a slow checkout.
|
||||
*/
|
||||
@Test
|
||||
void addCleansUpWhenTheWorktreeAddRunnerFailsAfterCreatingState(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
String branch = "cb-309-timeout";
|
||||
WorktreeException timeout = new WorktreeException("command timed out: synthetic git worktree add");
|
||||
AtomicReference<String> createdPath = new AtomicReference<>();
|
||||
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString(), null, null, null, command -> {
|
||||
createdPath.set(command[5]);
|
||||
try {
|
||||
git(repo, "worktree", "add", command[5], "-b", command[7], command[8]);
|
||||
} catch (Exception e) {
|
||||
throw new AssertionError("test setup could not create the worktree", e);
|
||||
}
|
||||
throw timeout;
|
||||
});
|
||||
|
||||
WorktreeException thrown = assertThrows(WorktreeException.class,
|
||||
() -> worktrees.add(repo.toString(), branch, "HEAD"));
|
||||
|
||||
assertSame(timeout, thrown, "cleanup must not replace the add failure");
|
||||
assertNotNull(createdPath.get(), "the add runner must receive the worktree path");
|
||||
assertFalse(Files.exists(Path.of(createdPath.get())),
|
||||
"the worktree directory leaked after the add runner failed");
|
||||
assertFalse(refExists(repo, "refs/heads/" + branch), "the branch leaked after the add runner failed");
|
||||
}
|
||||
|
||||
/** An ordinary Git refusal must not delete the existing branch or log a cleanup warning. */
|
||||
@Test
|
||||
void addFailureBeforeCreatingAWorktreeIsQuiet(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
String branch = "already-exists";
|
||||
git(repo, "branch", branch);
|
||||
|
||||
assertThrows(WorktreeException.class, () -> new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), branch, "HEAD"));
|
||||
|
||||
assertTrue(refExists(repo, "refs/heads/" + branch), "the existing branch must remain");
|
||||
assertTrue(capturedMessages().isEmpty(), "an ordinary Git refusal logged a warning: " + capturedMessages());
|
||||
}
|
||||
|
||||
// ---- CB-189: broader remote-URL coverage — every remote, both fetch and push URLs, any
|
||||
// non-SSH scheme. Reporting only, additive to the origin/https strip-and-refuse tests above. ----
|
||||
|
||||
|
||||
@@ -26,7 +26,12 @@ import org.slf4j.LoggerFactory;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@@ -178,14 +183,19 @@ class SessionManagerTest {
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
|
||||
boolean clearAfterTurn) {
|
||||
boolean clearAfterTurn) {
|
||||
return sessionManager(herdr, clock, contextCap, clearAfterTurn, null);
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
|
||||
boolean clearAfterTurn, java.util.function.Consumer<MemberSession> hook) {
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn);
|
||||
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn, hook);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -670,6 +680,24 @@ class SessionManagerTest {
|
||||
"BUSY session remains");
|
||||
}
|
||||
|
||||
@Test
|
||||
void reapIdleDoesNotReleaseSessionDeliveredAfterItsEligibilityCheck() {
|
||||
long[] clock = {0};
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager[] manager = new SessionManager[1];
|
||||
SessionManager sessions = sessionManager(herdr, () -> clock[0], 0, false,
|
||||
session -> manager[0].onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId())));
|
||||
manager[0] = sessions;
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
|
||||
clock[0] = 11;
|
||||
assertEquals(0, sessions.reapIdle(10), "delivery replaces the idle snapshot before release");
|
||||
assertEquals(MemberSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"a just-delivered session stays registered and busy");
|
||||
assertFalse(herdr.called("pane.close"), "the busy session pane is not stopped");
|
||||
}
|
||||
|
||||
@Test
|
||||
void doneSessionPastIdleTtlIsReaped() {
|
||||
long[] clock = {0};
|
||||
@@ -824,6 +852,186 @@ class SessionManagerTest {
|
||||
.count();
|
||||
}
|
||||
|
||||
// --- fleetd #308: a spawn accepted while the shutdown drain is running must not orphan ---
|
||||
|
||||
@Test
|
||||
void acquireRefusesANewSpawnOnceDrainAllHasStarted() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
|
||||
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(50)); // empty roster — returns immediately,
|
||||
// but the shutdown guard it flips must stay tripped for the life of the process.
|
||||
|
||||
ShuttingDownException e = assertThrows(ShuttingDownException.class,
|
||||
() -> sessions.acquire("ltms-local", null, "/caller", "term_primary"),
|
||||
"a spawn requested after the drain has begun must be refused loudly (invariant 3), "
|
||||
+ "not silently registered into a registry the drain will never revisit");
|
||||
assertNotNull(e.getMessage());
|
||||
assertFalse(e.getMessage().isBlank(), "the refusal must say why, not just that it failed");
|
||||
assertTrue(sessions.roster().isEmpty(), "the refused spawn must never reach the registry");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #308: the guard above closes most of the shutdown-race window, but it cannot close
|
||||
* all of it — a caller that already passed the {@code draining} check before {@code drainAll}
|
||||
* flips it can still be mid-{@code launcher.spawn()} (a real herdr round trip, not
|
||||
* instantaneous) when {@code drainAll} takes its registry snapshot. This test forces exactly
|
||||
* that interleaving with a launcher double that blocks the second {@code spawn()} call and the
|
||||
* first {@code stop()} call until released, then proves the post-loop sweep in {@code
|
||||
* drainAll} still finds and tears down the straggler that lands in the registry afterward.
|
||||
*/
|
||||
@Test
|
||||
void drainAllSweepsAStragglerThatRegisteredAfterTheInitialSnapshot() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
RaceLauncher race = new RaceLauncher(delegate);
|
||||
SessionManager sessions = new SessionManager(race);
|
||||
|
||||
// Registered normally, before the drain starts — the first spawn call, never blocked.
|
||||
MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
|
||||
|
||||
ExecutorService exec = Executors.newFixedThreadPool(2);
|
||||
try {
|
||||
// The straggler's acquire() reads `draining == false` (checked before this call ever
|
||||
// touches the launcher) and then blocks inside its own spawn() — the second spawn call.
|
||||
Future<MemberSession> straggler = exec.submit(() ->
|
||||
sessions.acquire("ltms-local", "/late", "/caller", "ownerLate"));
|
||||
|
||||
assertTrue(race.enteredSecondSpawn.await(5, TimeUnit.SECONDS),
|
||||
"the straggler must have passed the shutdown guard and reached spawn() before "
|
||||
+ "drainAll ever runs");
|
||||
assertEquals(1, sessions.roster().size(),
|
||||
"the straggler is still inside spawn() — not registered yet");
|
||||
|
||||
// drainAll flips `draining`, snapshots the registry (only `ready` is in it), and starts
|
||||
// releasing that snapshot — its first release() call stops `ready`'s pane, which this
|
||||
// launcher double blocks on so the interleaving below is deterministic, not a timing bet.
|
||||
Future<?> drain = exec.submit(() -> sessions.drainAll(TimeUnit.SECONDS.toNanos(5)));
|
||||
|
||||
assertTrue(race.enteredFirstStop.await(5, TimeUnit.SECONDS),
|
||||
"drainAll must be stopping the ready session's pane — proof its initial "
|
||||
+ "registry snapshot has already been taken");
|
||||
|
||||
// Only now does the straggler's spawn complete and register — strictly after the
|
||||
// snapshot drainAll's main pass is working from.
|
||||
race.releaseSecondSpawn.countDown();
|
||||
MemberSession registered = straggler.get(5, TimeUnit.SECONDS);
|
||||
|
||||
// Let drainAll finish releasing `ready`; it then re-checks the registry and must find
|
||||
// (and drain) the straggler that just landed in it.
|
||||
race.releaseFirstStop.countDown();
|
||||
drain.get(5, TimeUnit.SECONDS);
|
||||
|
||||
assertTrue(sessions.roster().isEmpty(),
|
||||
"the post-loop sweep must drain the straggler too, not just the initial snapshot");
|
||||
assertNotNull(registered.paneId());
|
||||
long paneCloseCalls = herdr.calls.stream().filter(c -> "pane.close".equals(c.method())).count();
|
||||
assertEquals(2, paneCloseCalls,
|
||||
"both ready's pane AND the straggler's pane must actually be stopped — a pane "
|
||||
+ "left running is exactly the orphan this ticket is about");
|
||||
} finally {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delegates every call while blocking the SECOND {@code spawn()} call and the FIRST
|
||||
* {@code stop()} call until the test releases them — used to force the fleetd #308 race
|
||||
* deterministically instead of betting on real thread-scheduling timing.
|
||||
*/
|
||||
private static final class RaceLauncher implements PeerLauncher {
|
||||
private final PeerLauncher delegate;
|
||||
private final AtomicInteger spawnCalls = new AtomicInteger();
|
||||
private final AtomicInteger stopCalls = new AtomicInteger();
|
||||
final CountDownLatch enteredSecondSpawn = new CountDownLatch(1);
|
||||
final CountDownLatch releaseSecondSpawn = new CountDownLatch(1);
|
||||
final CountDownLatch enteredFirstStop = new CountDownLatch(1);
|
||||
final CountDownLatch releaseFirstStop = new CountDownLatch(1);
|
||||
|
||||
RaceLauncher(PeerLauncher delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
private static void awaitOrFail(CountDownLatch latch) {
|
||||
try {
|
||||
if (!latch.await(5, TimeUnit.SECONDS)) {
|
||||
throw new AssertionError("RaceLauncher latch timed out");
|
||||
}
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new AssertionError("RaceLauncher latch interrupted", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return delegate.capabilities();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilitiesFor(String profileName) {
|
||||
return delegate.capabilitiesFor(profileName);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
if (spawnCalls.incrementAndGet() == 2) {
|
||||
enteredSecondSpawn.countDown();
|
||||
awaitOrFail(releaseSecondSpawn);
|
||||
}
|
||||
return delegate.spawn(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return delegate.profiles();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return delegate.defaultProfile();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return delegate.effectiveCwd(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> parityOverlay(String profileName) {
|
||||
return delegate.parityOverlay(profileName);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<?> list() {
|
||||
return delegate.list();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() {
|
||||
return delegate.reapOrphanWorkers();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
if (stopCalls.incrementAndGet() == 1) {
|
||||
enteredFirstStop.countDown();
|
||||
awaitOrFail(releaseFirstStop);
|
||||
}
|
||||
delegate.stop(id);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean clearContext(String id) {
|
||||
return delegate.clearContext(id);
|
||||
}
|
||||
}
|
||||
|
||||
private static List<String> promptTexts(FakeHerdr herdr) {
|
||||
return herdr.calls.stream()
|
||||
.filter(c -> "agent.prompt".equals(c.method()))
|
||||
|
||||
Reference in New Issue
Block a user