f8522edacd
The bridge cannot classify what it never emits. Fleet-health monitoring — detect a wedged member, decide, escalate — is blocked on that, so this is the foundation rather than the feature. A survey of Injector, CompletionResolver, StatusPoller, SessionManager, SessionReaper and MessageService found six conditions that ended a member's usefulness while saying nothing useful: * Injector.drop() — worker gone, queue cleared: SILENT * CompletionResolver.fail() — "via turn-stall fallback" at DEBUG, no reason * SessionManager.onFailed() — "session marked failed" at DEBUG, no stage * SessionManager.reapIdle() — indistinguishable from any other release * SessionManager.acquire() — spawn failure rethrown with no log at all * MessageService.abandon() — failed a caller's request at DEBUG The first two are the exact phrases that misled the CB-560 diagnosis: both name a symptom and neither names a cause. They are now WARN and carry the reason, the stage, and the counts. Conditions already loud were left alone, and so were two by-design timeouts in MessageService — an async model exists precisely for those, and promoting them would turn healthy operation into noise. Behaviour is unchanged: every edit is a log statement. Merge note: the recycle test deleted by CB-565 conflicted with a test added here. Resolved by keeping the new onTurnFailed assertion and dropping the recycle test, which tests a method that no longer exists.
630 lines
29 KiB
Java
630 lines
29 KiB
Java
package dev.ltms.bridged.session;
|
|
|
|
import dev.ltms.bridged.auth.MemberLifecycle;
|
|
import dev.ltms.bridged.herdr.Agent;
|
|
import dev.ltms.bridged.inject.TurnListener;
|
|
import dev.ltms.bridged.inject.MemberPresence;
|
|
import dev.ltms.bridged.peer.MemberRole;
|
|
import dev.ltms.bridged.peer.PeerHandle;
|
|
import dev.ltms.bridged.peer.PeerLauncher;
|
|
import dev.ltms.bridged.peer.SpawnRequest;
|
|
import org.slf4j.Logger;
|
|
import org.slf4j.LoggerFactory;
|
|
|
|
import java.security.SecureRandom;
|
|
import java.util.LinkedHashMap;
|
|
import java.util.List;
|
|
import java.util.Map;
|
|
import java.util.Optional;
|
|
import java.util.concurrent.ConcurrentHashMap;
|
|
import java.util.concurrent.TimeUnit;
|
|
import java.util.concurrent.atomic.AtomicLong;
|
|
import java.util.function.Consumer;
|
|
import java.util.function.LongSupplier;
|
|
|
|
/**
|
|
* Authoritative in-daemon registry of the worker sessions this {@code bridged} process spawned.
|
|
* Delegates spawn/teardown to a {@link PeerLauncher} (which performs subscription-guarded env
|
|
* setup and process/materialization) and adds lifecycle tracking, ownership, and deterministic
|
|
* teardown on top.
|
|
*
|
|
* <p>The state machine is intentionally one-shot / no-reuse: every acquired worker is fresh,
|
|
* and a finished or released worker is torn down, never pooled or reused.
|
|
*
|
|
* <p>The manager implements {@link TurnListener} so the injector's turn boundaries drive
|
|
* {@code READY → BUSY → DONE} (or {@code FAILED}). It exposes a {@link MemberPresence} view via
|
|
* {@link #asPresence()}: any MCP contact from a worker marks it present and simultaneously
|
|
* transitions the session {@code SPAWNING → READY}.
|
|
*/
|
|
public final class SessionManager implements TurnListener {
|
|
|
|
private static final Logger log = LoggerFactory.getLogger(SessionManager.class);
|
|
|
|
private final PeerLauncher launcher;
|
|
private final Worktrees worktrees;
|
|
private final ConcurrentHashMap<String /*paneId*/, MemberSession> registry = new ConcurrentHashMap<>();
|
|
private final MemberPresence presence;
|
|
private final SecureRandom nonceRandom = new SecureRandom();
|
|
private final AtomicLong nonceSeq = new AtomicLong();
|
|
private final LongSupplier nowNanos;
|
|
private final int contextCap;
|
|
private final boolean clearAfterTurn;
|
|
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
|
|
|
|
/** 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 terminalId on every release; no-op until wired. */
|
|
private final List<Consumer<String>> releaseListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
|
|
|
|
/** Backward-compatible constructor: shared-tree sessions, production git seam. */
|
|
public SessionManager(PeerLauncher launcher) {
|
|
this(launcher, new GitWorktrees(), System::nanoTime, 0, false);
|
|
}
|
|
|
|
/** Backward-compatible constructor with an injectable worktree seam. */
|
|
public SessionManager(PeerLauncher launcher, Worktrees worktrees) {
|
|
this(launcher, worktrees, System::nanoTime, 0, false);
|
|
}
|
|
|
|
/** Test constructor with an injectable clock. */
|
|
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos) {
|
|
this(launcher, worktrees, nowNanos, 0, false);
|
|
}
|
|
|
|
/** Production constructor with a configured context turn cap. */
|
|
public SessionManager(PeerLauncher launcher, Worktrees worktrees, int contextCap) {
|
|
this(launcher, worktrees, System::nanoTime, contextCap, false);
|
|
}
|
|
|
|
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
|
|
int contextCap) {
|
|
this(launcher, worktrees, nowNanos, contextCap, false);
|
|
}
|
|
|
|
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
|
|
int contextCap, boolean clearAfterTurn) {
|
|
this.launcher = launcher;
|
|
this.worktrees = worktrees;
|
|
this.presence = new PresenceBridge(this);
|
|
this.nowNanos = nowNanos;
|
|
this.contextCap = contextCap;
|
|
this.clearAfterTurn = clearAfterTurn;
|
|
}
|
|
|
|
/**
|
|
* The single {@link MemberPresence} view of this manager: it records availability and forwards
|
|
* the signal to the {@code SPAWNING → READY} transition. Pass this to the {@code Injector} and
|
|
* {@code BridgeMcp} where they previously accepted a plain {@link MemberPresence}. The same
|
|
* instance is returned every call — presence is shared state, so a fresh bridge per call would
|
|
* fragment the {@code present} set and lose signals across callers.
|
|
*/
|
|
public MemberPresence asPresence() {
|
|
return presence;
|
|
}
|
|
|
|
/**
|
|
* Spawn a worker and register it as {@link MemberSession.State#SPAWNING}. The caller's
|
|
* identity is recorded as {@code ownerTerminal} ({@code null} for daemon/anon callers).
|
|
*/
|
|
public MemberSession acquire(String profile, String requestedCwd, String callerCwd,
|
|
String ownerTerminal) {
|
|
return acquire(profile, requestedCwd, callerCwd, ownerTerminal, null);
|
|
}
|
|
|
|
/**
|
|
* Spawn a member, optionally inside a fresh git worktree, defaulting the role to
|
|
* {@link MemberRole#DEV}.
|
|
*
|
|
* <p>{@code DEV} is the right default because it is exactly what the old "worker" meant: an
|
|
* unqualified spawn is a unit of implementation work. An architect or a reviewer is always
|
|
* asked for on purpose, so neither is ever what a caller silently gets.
|
|
*/
|
|
public MemberSession acquire(String profile, String requestedCwd, String callerCwd,
|
|
String ownerTerminal, WorktreeRequest wt) {
|
|
return acquire(profile, MemberRole.DEV, requestedCwd, callerCwd, ownerTerminal, wt);
|
|
}
|
|
|
|
/**
|
|
* Spawn a member, optionally inside a fresh git worktree. When {@code wt} is non-null the
|
|
* worktree is provisioned, parity-overlaid, and its path becomes the member's cwd. On any
|
|
* failure before registration the worktree is removed so no dangling checkout is left.
|
|
*
|
|
* @param profile which backend to run on — a {@code profiles:} key
|
|
* @param role which contract the member runs under; never {@code null}
|
|
*/
|
|
public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd,
|
|
String ownerTerminal, WorktreeRequest wt) {
|
|
MemberRole memberRole = (role == null) ? MemberRole.DEV : role;
|
|
if (wt == null) {
|
|
// CB-557: the role must ride on the SpawnRequest, not stay a local. The launcher needs it
|
|
// to pick the profile out of that role's pool and to label the tab; a role kept only on
|
|
// the MemberSession is recorded after the spawn it was supposed to steer.
|
|
SpawnRequest req = new SpawnRequest(profile, requestedCwd, callerCwd, null, null, memberRole);
|
|
PeerHandle handle;
|
|
try {
|
|
handle = launcher.spawn(req);
|
|
} catch (RuntimeException e) {
|
|
log.warn("spawn failed for profile={} role={}: {}", profile, memberRole, e.getMessage());
|
|
throw e;
|
|
}
|
|
String resolvedProfile = resolveProfile(handle, profile);
|
|
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd));
|
|
long now = nowNanos.getAsLong();
|
|
MemberSession session = new MemberSession(
|
|
handle.id(),
|
|
handle.terminalId(),
|
|
resolvedProfile,
|
|
memberRole,
|
|
cwd,
|
|
ownerTerminal,
|
|
now,
|
|
now,
|
|
0,
|
|
MemberSession.State.SPAWNING,
|
|
null,
|
|
null);
|
|
registry.put(handle.id(), session);
|
|
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
|
log.debug("acquired session id={} terminal={} profile={} owner={}",
|
|
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
|
|
notifyAcquired(session.terminalId());
|
|
return session;
|
|
}
|
|
return acquireWithWorktree(profile, memberRole, requestedCwd, callerCwd, ownerTerminal, wt);
|
|
}
|
|
|
|
/** Tear a worker down by pane id and remove it from the registry. Idempotent. */
|
|
public void release(String paneId) {
|
|
release(paneId, ReleaseCause.COMPLETED);
|
|
}
|
|
|
|
/**
|
|
* Core teardown: always stops the worker pane and deregisters the session; whether the worker's
|
|
* git worktree is also removed depends on {@code cause}.
|
|
*
|
|
* <p>CB-544: these are two concerns that used to be fused. Stopping the pane is correct on every
|
|
* teardown — the worker process must end. Removing the worktree is a destructive act that is only
|
|
* correct for a deliberately-finished teardown (an explicit stop of a completed session, the
|
|
* reaper releasing a genuinely idle one, or a context-capped session). A shutdown drain
|
|
* must stop panes but preserve worktrees: a worker's uncommitted work exists in exactly one
|
|
* place — its worktree — so deleting it while the daemon simply goes down is silent data loss,
|
|
* with no copy and no error. Do NOT fuse these back together; the cost of an orphaned worktree
|
|
* is a logged path an operator can reclaim, the cost of a deleted one is unrecoverable work.
|
|
*/
|
|
private void release(String paneId, ReleaseCause cause) {
|
|
MemberSession removed = registry.remove(paneId);
|
|
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
|
|
if (removed != null) {
|
|
memberLifecycle.released(removed.terminalId());
|
|
log.debug("releasing session pane={} terminal={} state={} cause={}",
|
|
removed.paneId(), removed.terminalId(), removed.state(), cause);
|
|
if (preserveWorktree && removed.worktree() != null) {
|
|
logPreservedForShutdown(removed);
|
|
}
|
|
// CB-516: a send still waiting on this worker can never be answered now. Tell the
|
|
// listener BEFORE the pane is torn down, so a blocked caller fails fast with a real
|
|
// reason instead of sitting on a rendezvous nothing will ever resolve.
|
|
notifyReleased(removed.terminalId());
|
|
}
|
|
launcher.stop(paneId);
|
|
if (removed != null && !preserveWorktree && removed.worktree() != null) {
|
|
worktrees.remove(worktrees.repoRoot(removed.cwd()), removed.worktree());
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Why a session is being released — governs whether its worktree is preserved or removed.
|
|
* Worktree removal is reserved for the one case that is genuinely finished; everything else
|
|
* must keep the worker's only copy of its work.
|
|
*/
|
|
public enum ReleaseCause {
|
|
/** Deliberate teardown of a finished session. Stops the pane and removes the worktree. */
|
|
COMPLETED,
|
|
/** Daemon shutdown drain. Stops the pane but PRESERVES the worktree. */
|
|
SHUTDOWN
|
|
}
|
|
|
|
/**
|
|
* CB-544 shutdown drain log for a worktree we deliberately kept. A session still {@code BUSY}
|
|
* when the drain timeout expired was abandoned mid-turn — that work may be uncommitted and is
|
|
* the only copy — so the message is loud and points at the path an operator needs to reclaim.
|
|
*/
|
|
private void logPreservedForShutdown(MemberSession session) {
|
|
if (session.state() == MemberSession.State.BUSY) {
|
|
log.warn("shutdown drain abandoned BUSY session pane={} terminal={} mid-turn; "
|
|
+ "worktree preserved at {}", session.paneId(), session.terminalId(),
|
|
session.worktree());
|
|
} else {
|
|
log.info("shutdown drain preserved worktree at {} for pane={}",
|
|
session.worktree(), session.paneId());
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Register a callback invoked with a session's {@code terminalId} whenever it is acquired
|
|
* (CB-520). This is the hook that lets the reply inbox {@code own} a target's queue.
|
|
*/
|
|
public void onAcquire(Consumer<String> listener) {
|
|
if (listener != null) {
|
|
acquireListeners.add(listener);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Register a callback invoked with a session's {@code terminalId} whenever it is released
|
|
* (CB-516). Every teardown path funnels through {@link #release}, so one hook covers the REST
|
|
* and MCP stop tools, the idle-TTL reaper, and shutdown drain alike.
|
|
*
|
|
* <p>Added rather than injected because {@code MessageService} — one intended listener — is
|
|
* constructed after this manager (it needs the injector and rendezvous, which need the session
|
|
* presence view this manager exposes). Wiring it at construction would require breaking that
|
|
* cycle for one callback.
|
|
*/
|
|
public void onRelease(Consumer<String> listener) {
|
|
if (listener != null) {
|
|
releaseListeners.add(listener);
|
|
}
|
|
}
|
|
|
|
/** Inject the optional member-slot lifecycle after construction without changing constructors. */
|
|
public void setMemberLifecycle(MemberLifecycle memberLifecycle) {
|
|
this.memberLifecycle = memberLifecycle == null ? MemberLifecycle.NONE : memberLifecycle;
|
|
}
|
|
|
|
/** A listener failure must never prevent the acquisition it is reacting to. */
|
|
private void notifyAcquired(String terminalId) {
|
|
if (terminalId == null) {
|
|
return;
|
|
}
|
|
for (Consumer<String> listener : acquireListeners) {
|
|
try {
|
|
listener.accept(terminalId);
|
|
} catch (RuntimeException e) {
|
|
log.warn("acquire listener failed for terminal {}: {}", terminalId, e.toString());
|
|
}
|
|
}
|
|
}
|
|
|
|
/** A listener failure must never prevent the teardown it is reacting to. */
|
|
private void notifyReleased(String terminalId) {
|
|
if (terminalId == null) {
|
|
return;
|
|
}
|
|
for (Consumer<String> listener : releaseListeners) {
|
|
try {
|
|
listener.accept(terminalId);
|
|
} catch (RuntimeException e) {
|
|
log.warn("release listener failed for terminal {}: {}", terminalId, e.toString());
|
|
}
|
|
}
|
|
}
|
|
|
|
private MemberSession acquireWithWorktree(String profile, MemberRole memberRole, String requestedCwd, String callerCwd,
|
|
String ownerTerminal, WorktreeRequest wt) {
|
|
String preResolvedProfile = (profile == null || profile.isBlank())
|
|
? launcher.defaultProfile() : profile;
|
|
// CB-507: resolve through the launcher's CB-112 chain (requested → profile cwd → caller →
|
|
// daemon cwd → "."), never the raw args. A plain REST spawn supplies neither a requested
|
|
// nor a caller cwd, so taking the first non-blank of those two yielded null and put
|
|
// `git -C null` on the command line — an NPE out of ProcessBuilder, surfacing as HTTP 500.
|
|
// The non-worktree path always used this chain; only this branch was missed.
|
|
String repoRoot = worktrees.repoRoot(
|
|
launcher.effectiveCwd(new SpawnRequest(preResolvedProfile, requestedCwd, callerCwd)));
|
|
String branch = "worker/" + slug(wt.ticketSlug()) + "-" + nonce();
|
|
String path = null;
|
|
PeerHandle handle;
|
|
try {
|
|
path = worktrees.add(repoRoot, branch, wt.baseRef());
|
|
worktrees.overlayParity(repoRoot, path, launcher.parityOverlay(preResolvedProfile));
|
|
handle = launcher.spawn(new SpawnRequest(profile, path, callerCwd, null, null, memberRole));
|
|
} catch (RuntimeException e) {
|
|
log.warn("spawn failed for profile={} role={} branch={} path={}: {}",
|
|
preResolvedProfile, memberRole, branch, path, e.getMessage());
|
|
if (path != null) {
|
|
try {
|
|
worktrees.remove(repoRoot, path);
|
|
} catch (RuntimeException cleanup) {
|
|
log.warn("failed to clean up worktree {} after spawn error: {}", path, cleanup.getMessage());
|
|
}
|
|
}
|
|
throw e;
|
|
}
|
|
String resolvedProfile = resolveProfile(handle, profile);
|
|
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, path, callerCwd));
|
|
long now = nowNanos.getAsLong();
|
|
MemberSession session = new MemberSession(
|
|
handle.id(),
|
|
handle.terminalId(),
|
|
resolvedProfile,
|
|
memberRole,
|
|
cwd,
|
|
ownerTerminal,
|
|
now,
|
|
now,
|
|
0,
|
|
MemberSession.State.SPAWNING,
|
|
path,
|
|
branch);
|
|
registry.put(handle.id(), session);
|
|
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
|
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
|
|
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
|
|
notifyAcquired(session.terminalId());
|
|
return session;
|
|
}
|
|
|
|
private String slug(String raw) {
|
|
return raw == null ? "ticket" : raw.toLowerCase().replaceAll("[^a-z0-9]+", "-").replaceAll("^-+|-+$", "");
|
|
}
|
|
|
|
private String nonce() {
|
|
return String.format("%06x", nonceRandom.nextInt(1 << 24)) + "-" + nonceSeq.incrementAndGet();
|
|
}
|
|
|
|
/**
|
|
* The profile to record for a session. A launcher that performed dynamic selection tells us
|
|
* the actual profile via {@link PeerHandle#profile()}; otherwise fall back to what the caller
|
|
* requested (or the launcher's default for a no-profile spawn).
|
|
*/
|
|
private String resolveProfile(PeerHandle handle, String requestedProfile) {
|
|
String fromHandle = handle.profile();
|
|
if (fromHandle != null && !fromHandle.isBlank()) {
|
|
return fromHandle;
|
|
}
|
|
if (requestedProfile != null && !requestedProfile.isBlank()) {
|
|
return requestedProfile;
|
|
}
|
|
return launcher.defaultProfile();
|
|
}
|
|
|
|
/** The session for {@code paneId}, if it is still registered and not released. */
|
|
public Optional<MemberSession> get(String paneId) {
|
|
return Optional.ofNullable(registry.get(paneId));
|
|
}
|
|
|
|
/** Bridge-owned roster: all registered sessions (acquired minus released). */
|
|
public List<MemberSession> roster() {
|
|
return List.copyOf(registry.values());
|
|
}
|
|
|
|
/**
|
|
* CB-304 merged roster+live view. The registry is authoritative for worktree, branch,
|
|
* profile, owner, and state; the optional live agent supplies the herdr-reported status.
|
|
*/
|
|
public static Map<String, Object> rosterView(MemberSession session, Agent live) {
|
|
Map<String, Object> m = new LinkedHashMap<>();
|
|
m.put("sessionId", session.terminalId());
|
|
m.put("paneId", session.paneId());
|
|
// Both axes, always: profile says which backend this member runs on, role says what it is
|
|
// for. A lead reading the roster needs both — two rows may share a profile and still be
|
|
// allowed to do entirely different things.
|
|
m.put("profile", session.profile());
|
|
m.put("role", session.role() == null ? "dev" : session.role().wireName());
|
|
m.put("state", session.state().name().toLowerCase());
|
|
if (session.worktree() != null) {
|
|
m.put("worktree", session.worktree());
|
|
}
|
|
if (session.branch() != null) {
|
|
m.put("branch", session.branch());
|
|
}
|
|
if (session.ownerTerminal() != null) {
|
|
m.put("owner", session.ownerTerminal());
|
|
}
|
|
m.put("liveStatus", live == null ? "unknown" : live.status().name().toLowerCase());
|
|
return m;
|
|
}
|
|
|
|
/** Lifecycle hook: worker became available on the bridge MCP. */
|
|
void onReady(String terminalId) {
|
|
transitionByTerminal(terminalId, MemberSession.State.SPAWNING, MemberSession.State.READY);
|
|
}
|
|
|
|
/**
|
|
* Lifecycle hook: a message was delivered into the worker — it is now busy on a turn.
|
|
* The turn count is bumped and the activity timestamp is refreshed. A {@code DONE} session
|
|
* can be re-delivered for multi-turn reuse until it is released.
|
|
*/
|
|
@Override
|
|
public void onDelivered(String target) {
|
|
MemberSession current = findByTerminal(target);
|
|
if (current == null) return;
|
|
if (current.state() != MemberSession.State.READY && current.state() != MemberSession.State.DONE) {
|
|
return;
|
|
}
|
|
long now = nowNanos.getAsLong();
|
|
MemberSession updated = current.withState(MemberSession.State.BUSY).bumpTurn(now);
|
|
if (replace(current, updated)) {
|
|
log.debug("session transitioned terminal={} pane={} {} -> BUSY turn={}",
|
|
target, current.paneId(), current.state(), updated.turnCount());
|
|
}
|
|
}
|
|
|
|
/** Lifecycle hook: the worker's delegated turn completed successfully. */
|
|
@Override
|
|
public void onTurnComplete(String target) {
|
|
completeTurn(target, false);
|
|
}
|
|
|
|
@Override
|
|
public boolean hasPostTurnAction(String target) {
|
|
if (!clearAfterTurn) return false;
|
|
MemberSession current = findByTerminal(target);
|
|
return current != null && current.state() == MemberSession.State.BUSY
|
|
&& (contextCap <= 0 || current.turnCount() < contextCap);
|
|
}
|
|
|
|
@Override
|
|
public boolean onTurnCompleteWithPostAction(String target) {
|
|
return completeTurn(target, true);
|
|
}
|
|
|
|
private boolean completeTurn(String target, boolean startContextReset) {
|
|
MemberSession current = findByTerminal(target);
|
|
if (current == null || current.state() != MemberSession.State.BUSY) return false;
|
|
long now = nowNanos.getAsLong();
|
|
MemberSession updated = current.withState(MemberSession.State.DONE).withActivity(now);
|
|
if (replace(current, updated)) {
|
|
log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}",
|
|
target, current.paneId(), updated.turnCount());
|
|
}
|
|
if (contextCap > 0 && updated.turnCount() >= contextCap) {
|
|
release(current.paneId());
|
|
return false;
|
|
}
|
|
if (!startContextReset || !clearAfterTurn) return false;
|
|
try {
|
|
return launcher.clearContext(current.paneId());
|
|
} catch (RuntimeException e) {
|
|
log.warn("context reset failed for terminal={} pane={}; continuing without reset: {}",
|
|
target, current.paneId(), e.getMessage());
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/** Lifecycle hook: the worker's delegated turn failed. */
|
|
@Override
|
|
public void onTurnFailed(String target) {
|
|
onFailed(target);
|
|
}
|
|
|
|
/** Lifecycle hook: the worker vanished or was dropped mid-life. */
|
|
void onFailed(String target) {
|
|
MemberSession current = findByTerminal(target);
|
|
if (current == null) return;
|
|
if (current.state() == MemberSession.State.RELEASED) return;
|
|
MemberSession.State priorState = current.state();
|
|
if (replace(current, current.withState(MemberSession.State.FAILED))) {
|
|
log.warn("member terminal={} pane={} can no longer be delegated to: its turn never resolved "
|
|
+ "(was {} when it failed)", target, current.paneId(), priorState);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Best-effort reap of sessions that have been idle longer than {@code idleTtlNanos}. Only
|
|
* {@code READY} and {@code DONE} sessions are eligible — never a {@code SPAWNING} or
|
|
* {@code BUSY} worker. Returns the number of sessions released.
|
|
*/
|
|
int reapIdle(long idleTtlNanos) {
|
|
long now = nowNanos.getAsLong();
|
|
int reaped = 0;
|
|
for (MemberSession s : roster()) {
|
|
if (s.state() != MemberSession.State.READY && s.state() != MemberSession.State.DONE) {
|
|
continue;
|
|
}
|
|
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));
|
|
release(s.paneId());
|
|
reaped++;
|
|
}
|
|
}
|
|
return reaped;
|
|
}
|
|
|
|
/**
|
|
* Gracefully drain all registered sessions on daemon shutdown. For each session that is
|
|
* {@code BUSY}, poll up to {@code timeoutNanos} for it to leave {@code BUSY}, then release it
|
|
* regardless. Non-busy sessions are released immediately. A failure releasing one session is
|
|
* logged and does not abort the rest.
|
|
*
|
|
* <p>CB-544: this is a {@link ReleaseCause#SHUTDOWN} release — the worker's pane is stopped
|
|
* (the process must end) but its worktree is preserved and its path logged. Shutdown is never
|
|
* 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.
|
|
*/
|
|
void drainAll(long timeoutNanos) {
|
|
long deadline = System.nanoTime() + timeoutNanos;
|
|
for (MemberSession s : roster()) {
|
|
try {
|
|
if (s.state() == MemberSession.State.BUSY) {
|
|
while (System.nanoTime() < deadline) {
|
|
MemberSession current = registry.get(s.paneId());
|
|
if (current == null || current.state() != MemberSession.State.BUSY) {
|
|
break;
|
|
}
|
|
try {
|
|
long remaining = deadline - System.nanoTime();
|
|
Thread.sleep(Math.min(TimeUnit.NANOSECONDS.toMillis(remaining), 50));
|
|
} catch (InterruptedException e) {
|
|
Thread.currentThread().interrupt();
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
release(s.paneId(), ReleaseCause.SHUTDOWN);
|
|
} catch (RuntimeException e) {
|
|
log.warn("drain failed for pane={}; continuing with remaining sessions", s.paneId(), e);
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Close this manager by draining all sessions. The timeout comes from configuration when set,
|
|
* otherwise a sensible default.
|
|
*/
|
|
public void close(Integer drainTimeoutSeconds) {
|
|
int seconds = (drainTimeoutSeconds != null && drainTimeoutSeconds > 0) ? drainTimeoutSeconds : 5;
|
|
drainAll(TimeUnit.SECONDS.toNanos(seconds));
|
|
}
|
|
|
|
/** Number of sessions currently registered. */
|
|
public int size() {
|
|
return registry.size();
|
|
}
|
|
|
|
/**
|
|
* The registered session owning {@code terminalId}, or {@code null} if none does.
|
|
*
|
|
* <p>A null {@code terminalId} is a normal input, not a caller bug: every lifecycle hook here is
|
|
* fed from the MCP transport, where the <em>primary</em> resolves to a {@link
|
|
* dev.ltms.bridged.auth.Principal} with no terminal. {@code BridgeMcp} documents that contact as
|
|
* a no-op, and {@link dev.ltms.bridged.inject.MemberPresence#markPresent} honours it — but
|
|
* {@code PresenceBridge} then forwards the same null here. Matching on a null id can never
|
|
* succeed anyway (a registered session always has a terminal), so answer "no match" rather than
|
|
* throwing: an NPE on this path takes down an unrelated tool call for the primary.
|
|
*/
|
|
private MemberSession findByTerminal(String terminalId) {
|
|
if (terminalId == null) return null;
|
|
for (MemberSession s : registry.values()) {
|
|
if (terminalId.equals(s.terminalId())) return s;
|
|
}
|
|
return null;
|
|
}
|
|
|
|
private void transitionByTerminal(String terminalId, MemberSession.State from,
|
|
MemberSession.State to) {
|
|
MemberSession current = findByTerminal(terminalId);
|
|
if (current == null || current.state() != from) return;
|
|
long now = nowNanos.getAsLong();
|
|
if (replace(current, current.withState(to).withActivity(now))) {
|
|
log.debug("session transitioned terminal={} pane={} {} -> {}",
|
|
terminalId, current.paneId(), from, to);
|
|
}
|
|
}
|
|
|
|
private boolean replace(MemberSession expected, MemberSession updated) {
|
|
return registry.replace(expected.paneId(), expected, updated);
|
|
}
|
|
|
|
/** MemberPresence bridge that also drives the manager's READY transition. */
|
|
private static final class PresenceBridge extends MemberPresence {
|
|
private final SessionManager sessions;
|
|
|
|
PresenceBridge(SessionManager sessions) {
|
|
this.sessions = sessions;
|
|
}
|
|
|
|
@Override
|
|
public void markPresent(String terminal) {
|
|
if (terminal == null || terminal.isBlank()) {
|
|
return; // the primary's contact carries no worker terminal — not a readiness signal
|
|
}
|
|
super.markPresent(terminal);
|
|
sessions.onReady(terminal);
|
|
}
|
|
}
|
|
}
|