15ff6bcde5
A refs/wip snapshot ref is deleted only when both hold: its commit's tree is
already reachable from main, and it is older than 24h. Reachability is the safety
floor — a snapshot exists because the work was committed nowhere else, so an
unreachable one is the last copy and is never swept. Every deletion logs the ref
and the sha.
/members gains wipRefs{count,costBytes} so the growth is visible.
Verified against real git, not only the fakes: a recoverable+old ref is deleted,
a recoverable+young one survives the age floor, and the last copy survives. A repo
with no main deletes nothing, and a repo with no snapshots is a clean no-op.
Closes #67
831 lines
41 KiB
Java
831 lines
41 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.msg.TurnToken;
|
|
import dev.ltms.bridged.peer.Capability;
|
|
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.Set;
|
|
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-586: the repo root the fleet actually works in, remembered the first time a worktree
|
|
* session is spawned (worktrees are checkouts of it). {@code refs/wip/*} live there, and this
|
|
* single cached value is what the snapshot retention sweep and the operator-visible census run
|
|
* against. The daemon is bridged into one project at a time, so "the first worktree's repo" is
|
|
* the repo; {@code null} until any worktree is spawned, meaning nothing to sweep or measure.
|
|
*/
|
|
private volatile String fleetRepoRoot;
|
|
|
|
/** 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. */
|
|
private final List<Consumer<ReleaseDetail>> 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, with no session identity requested.
|
|
* Equivalent to {@link #acquire(String, MemberRole, String, String, String, WorktreeRequest,
|
|
* String, String)} with both trailing args {@code null}.
|
|
*
|
|
* @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) {
|
|
return acquire(profile, role, requestedCwd, callerCwd, ownerTerminal, wt, null, null);
|
|
}
|
|
|
|
/**
|
|
* Spawn a member, optionally inside a fresh git worktree, optionally onto a chosen or resumed
|
|
* agent session (CB-547a / CB-584). 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.
|
|
*
|
|
* <p>A non-blank {@code resumeSessionId} requires an explicit {@code profile}: a resumed
|
|
* conversation is tied to the specific backend that started it, so an unqualified spawn (whose
|
|
* backend a placement policy picks at spawn time) has no safe candidate to check the capability
|
|
* against. It also requires that profile's adapter to declare {@link Capability#SESSION_RESUME};
|
|
* refusing rather than silently delivering a cold session on a non-supporting adapter is the
|
|
* whole point of checking before spawn, not after (CB-584).
|
|
*
|
|
* @param profile which backend to run on — a {@code profiles:} key
|
|
* @param role which contract the member runs under; never {@code null}
|
|
* @param sessionName the bridge's logical name for the session, or {@code null}
|
|
* @param resumeSessionId the prior agent session to resume, or {@code null} for a fresh one
|
|
* @throws IllegalArgumentException if {@code resumeSessionId} is set with no explicit profile,
|
|
* or the resolved profile's adapter lacks
|
|
* {@link Capability#SESSION_RESUME}
|
|
*/
|
|
public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd,
|
|
String ownerTerminal, WorktreeRequest wt,
|
|
String sessionName, String resumeSessionId) {
|
|
MemberRole memberRole = (role == null) ? MemberRole.DEV : role;
|
|
requireResumeCapability(profile, resumeSessionId);
|
|
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, sessionName, resumeSessionId, 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,
|
|
handle.charterReceipt(),
|
|
handle.agentSessionId());
|
|
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,
|
|
sessionName, resumeSessionId);
|
|
}
|
|
|
|
/**
|
|
* Refuse a {@code resumeSessionId} that either names no explicit profile or names one whose
|
|
* adapter does not declare {@link Capability#SESSION_RESUME}. A no-op when
|
|
* {@code resumeSessionId} is blank — the ordinary, no-identity spawn path.
|
|
*/
|
|
private void requireResumeCapability(String profile, String resumeSessionId) {
|
|
if (resumeSessionId == null || resumeSessionId.isBlank()) {
|
|
return;
|
|
}
|
|
if (profile == null || profile.isBlank()) {
|
|
throw new IllegalArgumentException("resumeSessionId requires an explicit profile — "
|
|
+ "a resumed conversation is tied to the backend that started it, so it cannot "
|
|
+ "be left to placement to pick");
|
|
}
|
|
Set<Capability> caps = launcher.capabilitiesFor(profile);
|
|
if (!caps.contains(Capability.SESSION_RESUME)) {
|
|
throw new IllegalArgumentException("worker profile '" + profile + "' does not declare "
|
|
+ "Capability.SESSION_RESUME — refusing resumeSessionId rather than silently "
|
|
+ "starting a cold session");
|
|
}
|
|
}
|
|
|
|
/** 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;
|
|
String snapshotRef = null;
|
|
if (removed != null) {
|
|
try {
|
|
memberLifecycle.released(removed.terminalId());
|
|
log.debug("releasing session pane={} terminal={} state={} cause={}",
|
|
removed.paneId(), removed.terminalId(), removed.state(), cause);
|
|
boolean dirty = removed.worktree() != null && worktrees.hasUncommitted(removed.worktree());
|
|
if (preserveWorktree && removed.worktree() != null) {
|
|
logPreservedForShutdown(removed);
|
|
} else if (dirty) {
|
|
// CB-576: a release that would otherwise remove the worktree finds it holding
|
|
// uncommitted work the bridge cannot see. A worker that ends a turn without
|
|
// committing (normally because it stopped to ask a question or refused the turn)
|
|
// has its only copy of that work in the worktree. Remove would --force-delete it,
|
|
// so preserve the directory and tell an operator where to find it.
|
|
preserveWorktree = true;
|
|
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
|
|
+ "the worktree holds uncommitted changes that --force remove would destroy",
|
|
cause, removed.worktree(), removed.paneId(), removed.terminalId());
|
|
}
|
|
if (dirty) {
|
|
// CB-578 stage C: preserving on disk is not saving — the directory is one
|
|
// `worktree remove --force`, or an operator tidying up, away from gone. Commit
|
|
// its full state to a ref before the preserve-or-remove decision above can be
|
|
// undone by anything else, regardless of why this release fired.
|
|
snapshotRef = trySnapshot(removed, cause);
|
|
}
|
|
} catch (RuntimeException e) {
|
|
// CB-581: hasUncommitted shells out to `git status` and can throw on a non-zero
|
|
// exit. We can no longer tell whether the worktree holds uncommitted work, so fail
|
|
// toward the safe answer and preserve it — deleting on a guess can destroy work
|
|
// that has no other copy (CB-576), while keeping it on a false alarm only costs
|
|
// disk. The exception must not propagate: the pane still has to stop below.
|
|
preserveWorktree = true;
|
|
log.warn("release {} could not tell whether worktree {} for pane={} terminal={} has "
|
|
+ "uncommitted changes; preserving it rather than risk destroying unsaved work: {}",
|
|
cause, removed.worktree(), removed.paneId(), removed.terminalId(), e.toString());
|
|
} finally {
|
|
// CB-516/CB-581: a send still waiting on this worker can never be answered now, no
|
|
// matter what happened above. 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. CB-578 stage C: carry the worktree/branch/snapshot ref
|
|
// too, so a failed ticket's detail can point a lead at the same tree to re-dispatch.
|
|
// CB-584 (issue #65 criterion 5): carry agentSessionId alongside them, so a lead can
|
|
// also resume the member's conversation, not just re-dispatch onto its files.
|
|
notifyReleased(new ReleaseDetail(removed.terminalId(), removed.worktree(),
|
|
removed.branch(), snapshotRef, removed.agentSessionId()));
|
|
}
|
|
}
|
|
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
|
|
// 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) {
|
|
worktrees.remove(worktrees.repoRoot(removed.cwd()), removed.worktree());
|
|
}
|
|
}
|
|
|
|
/**
|
|
* 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
|
|
* regardless of whether this succeeds, so the only cost of a failed snapshot is a WARN and a
|
|
* missing ref — never a lost pane stop or a lost release notification.
|
|
*/
|
|
private String trySnapshot(MemberSession session, ReleaseCause cause) {
|
|
if (session.worktree() == null || session.branch() == null) {
|
|
return null;
|
|
}
|
|
try {
|
|
Optional<String> ref = worktrees.snapshot(session.worktree(), session.branch(),
|
|
snapshotMessage(session, cause));
|
|
ref.ifPresent(sha -> log.info(
|
|
"snapshotted dirty worktree {} to refs/wip/{} commit={} for pane={} terminal={}",
|
|
session.worktree(), session.branch(), sha, session.paneId(), session.terminalId()));
|
|
return ref.orElse(null);
|
|
} catch (RuntimeException e) {
|
|
log.warn("snapshot of dirty worktree {} failed for pane={} terminal={} branch={}: the "
|
|
+ "worktree is still preserved on disk, just not committed to refs/wip/{}: {}",
|
|
session.worktree(), session.paneId(), session.terminalId(), session.branch(),
|
|
session.branch(), e.toString());
|
|
return null;
|
|
}
|
|
}
|
|
|
|
/** Commit message for a CB-578 stage C snapshot — names the member so an operator can tell runs apart. */
|
|
private String snapshotMessage(MemberSession session, ReleaseCause cause) {
|
|
return "CB-578 stage C: snapshot of a released worker\n\n"
|
|
+ "terminal: " + session.terminalId() + "\n"
|
|
+ "profile: " + session.profile() + "\n"
|
|
+ "branch: " + session.branch() + "\n"
|
|
+ "cause: " + cause;
|
|
}
|
|
|
|
/**
|
|
* Facts about a released session that a listener needs beyond the bare terminal id — enough
|
|
* for a caller to point a lead at where to re-dispatch onto the same tree after a failed
|
|
* release (CB-578 stage C, acceptance criterion 10). {@code worktreePath} and {@code branch}
|
|
* are {@code null} for a shared-tree session; {@code snapshotRef} is {@code null} unless this
|
|
* release snapshotted a dirty worktree into {@code refs/wip/<branch>}.
|
|
*/
|
|
public record ReleaseDetail(String terminalId, String worktreePath, String branch, String snapshotRef,
|
|
String agentSessionId) {
|
|
}
|
|
|
|
/**
|
|
* 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<ReleaseDetail> 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(ReleaseDetail detail) {
|
|
if (detail.terminalId() == null) {
|
|
return;
|
|
}
|
|
for (Consumer<ReleaseDetail> listener : releaseListeners) {
|
|
try {
|
|
listener.accept(detail);
|
|
} catch (RuntimeException e) {
|
|
log.warn("release listener failed for terminal {}: {}", detail.terminalId(), e.toString());
|
|
}
|
|
}
|
|
}
|
|
|
|
private MemberSession acquireWithWorktree(String profile, MemberRole memberRole, String requestedCwd, String callerCwd,
|
|
String ownerTerminal, WorktreeRequest wt,
|
|
String sessionName, String resumeSessionId) {
|
|
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)));
|
|
if (fleetRepoRoot == null) {
|
|
// CB-586: remember the repo whose worktrees the fleet spawns — its refs/wip/* are the
|
|
// snapshot store the retention sweep and the operator census operate on.
|
|
fleetRepoRoot = repoRoot;
|
|
}
|
|
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, sessionName, resumeSessionId, 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,
|
|
handle.charterReceipt(),
|
|
handle.agentSessionId());
|
|
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());
|
|
}
|
|
// CB-584: which conversation this member holds — the id a later resumeSessionId spawn would
|
|
// pass back. Absent for an adapter that declines Capability.SESSION_RESUME, or one that
|
|
// resolves it lazily and has not yet.
|
|
if (session.agentSessionId() != null) {
|
|
m.put("agentSessionId", session.agentSessionId());
|
|
}
|
|
// CB-571: which charter this member was started with — never the charter text itself. The
|
|
// digest lets a lead tell at a glance whether all members got the same charter; the source
|
|
// records whether a role charter was configured ("fleet.charters.<role>") or only the reply
|
|
// charter was composed ("none").
|
|
if (session.charterReceipt() != null) {
|
|
m.put("charterSource", session.charterReceipt().charterSource());
|
|
if (session.charterReceipt().charterSha256() != null) {
|
|
m.put("charterSha256", session.charterReceipt().charterSha256());
|
|
}
|
|
}
|
|
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, TurnToken token) {
|
|
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));
|
|
// 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++;
|
|
} catch (RuntimeException e) {
|
|
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
|
|
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
|
|
}
|
|
}
|
|
}
|
|
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();
|
|
}
|
|
|
|
/**
|
|
* CB-586: the operator-visible census of {@code refs/wip/*} in the repo the fleet works in —
|
|
* how many snapshot refs exist and roughly what they cost. Empty (no repo known) until at
|
|
* least one worktree session has been spawned, exactly so a fleet that has never snapshotted
|
|
* anything surfaces nothing new, as it did before CB-586.
|
|
*/
|
|
public Optional<Worktrees.WipRefStats> wipRefs() {
|
|
String repo = fleetRepoRoot;
|
|
return repo == null ? Optional.empty() : Optional.of(worktrees.wipRefs(repo));
|
|
}
|
|
|
|
/**
|
|
* CB-586: run the snapshot retention sweep in the fleet's repo (a no-op until a worktree has
|
|
* been spawned, which establishes the repo). Returns how many {@code refs/wip/*} it deleted.
|
|
*/
|
|
public int sweepWipRefs(long minAgeMillis) {
|
|
String repo = fleetRepoRoot;
|
|
return repo == null ? 0 : worktrees.pruneWipRefs(repo, minAgeMillis);
|
|
}
|
|
|
|
/**
|
|
* 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);
|
|
}
|
|
}
|
|
}
|