diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index efff9dc..4b664e3 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -951,7 +951,10 @@ public final class FleetMcp { .sorted(Map.Entry.comparingByValue()) .map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm)) .toList(); - List roster = sessions.roster(); + // fleetd #209: this is the caller-driven fleet_list read that actually reports + // agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the + // resolving roster; the heartbeat/health/metrics timers stay on the plain sessions.roster(). + List roster = sessions.rosterResolved(); List> out = roster.stream() .map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong())) .toList(); diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index d78abf1..5cb89db 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -315,7 +315,9 @@ public final class FleetApp { .map(Agent.class::cast) .filter(a -> a.terminalId() != null) .collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b)); - List> out = sessions.roster().stream() + // fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it + // uses the resolving roster read (caller-driven, not a timer) rather than the plain one. + List> out = sessions.rosterResolved().stream() .map(s -> SessionManager.rosterView(s, live.get(s.terminalId()))) .toList(); Map body = new LinkedHashMap<>(); diff --git a/fleetd/src/main/java/dev/ltms/fleet/session/MemberSession.java b/fleetd/src/main/java/dev/ltms/fleet/session/MemberSession.java index 1c16201..f91fda4 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/session/MemberSession.java +++ b/fleetd/src/main/java/dev/ltms/fleet/session/MemberSession.java @@ -89,4 +89,15 @@ public record MemberSession( return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos, nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId); } + + /** + * Return a copy with {@code agentSessionId} resolved to a non-null value (fleetd #209). Some + * adapters (opencode) cannot answer {@link dev.ltms.fleet.peer.PeerHandle#agentSessionId()} at + * spawn time — the peer has not persisted its session record yet — so the id is discovered on + * a later poll and swapped into the otherwise-immutable session via this wither. + */ + public MemberSession withAgentSessionId(String agentSessionId) { + return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos, + lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId); + } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java b/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java index e135bb7..eb2f938 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java +++ b/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java @@ -46,6 +46,15 @@ public final class SessionManager implements TurnListener { private final PeerLauncher launcher; private final Worktrees worktrees; private final ConcurrentHashMap registry = new ConcurrentHashMap<>(); + /** + * fleetd #209: the live {@link PeerHandle} for every registered pane, retained solely so + * {@link #resolveAgentSessionId} can re-poll {@link PeerHandle#agentSessionId()} after spawn. + * The handle used to go out of scope at the end of the spawn method, so a launcher that answers + * the id lazily (opencode — the on-disk session row is written after the pane is created) could + * never be re-asked, and {@code fleet_list}/{@code fleet_spawn resumeSessionId} never saw it. + * Populated on every spawn path, removed on {@link #release}. + */ + private final ConcurrentHashMap handles = new ConcurrentHashMap<>(); private final MemberPresence presence; private final SecureRandom nonceRandom = new SecureRandom(); private final AtomicLong nonceSeq = new AtomicLong(); @@ -205,6 +214,7 @@ public final class SessionManager implements TurnListener { handle.charterReceipt(), handle.agentSessionId()); registry.put(handle.id(), session); + handles.put(handle.id(), handle); memberLifecycle.acquired(session.role(), session.profile(), session.terminalId()); log.debug("acquired session id={} terminal={} profile={} owner={}", handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal()); @@ -257,6 +267,10 @@ public final class SessionManager implements TurnListener { */ private void release(String paneId, ReleaseCause cause) { MemberSession removed = registry.remove(paneId); + // 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) { @@ -303,8 +317,11 @@ public final class SessionManager implements TurnListener { // 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())); + // fleetd #209: a late-resolving adapter (opencode) may only now have an id — resolve + // one last time so a released member's detail carries the id it now has. + MemberSession resolved = resolveAgentSessionId(removed, removedHandle); + notifyReleased(new ReleaseDetail(resolved.terminalId(), resolved.worktree(), + resolved.branch(), snapshotRef, resolved.agentSessionId())); } } // CB-581: the pane must always stop, even if the dirty check above threw. A session removed @@ -508,6 +525,7 @@ public final class SessionManager implements TurnListener { handle.charterReceipt(), handle.agentSessionId()); registry.put(handle.id(), session); + handles.put(handle.id(), handle); 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()); @@ -539,16 +557,73 @@ public final class SessionManager implements TurnListener { return launcher.defaultProfile(); } - /** The session for {@code paneId}, if it is still registered and not released. */ + /** + * The session for {@code paneId}, if it is still registered and not released. fleetd #209: + * resolves a still-unknown {@code agentSessionId} against the retained handle before returning, + * so {@code fleet_status} sees an id a lazy-resolving adapter has since written. + */ public Optional get(String paneId) { - return Optional.ofNullable(registry.get(paneId)); + return Optional.ofNullable(registry.get(paneId)).map(this::resolveAgentSessionId); } - /** Fleet-owned roster: all registered sessions (acquired minus released). */ + /** + * Fleet-owned roster: all registered sessions (acquired minus released). Deliberately does + * not resolve {@code agentSessionId} (fleetd #209 follow-up) — this is the + * roster supplier on the heartbeat and health-tick timers ({@code LeadHeartbeatLoop}, + * {@code FleetHealthMonitor} in {@code Fleetd}), on the placement/exhaustion paths, and on the + * metrics scrape ({@code FleetMetrics}), all called far more often than any caller actually + * reads {@code agentSessionId}. Resolving here would mean every tick opens a lazy-resolving + * adapter's (opencode's) on-disk session store once per member whose id is still unknown — and + * for a member whose id never appears, that cost never stops, for the life of the process. Use + * {@link #rosterResolved()} instead wherever the id must be current. + */ public List roster() { return List.copyOf(registry.values()); } + /** + * {@link #roster()}, with each session's still-unknown {@code agentSessionId} re-resolved + * against its retained handle (fleetd #209) — so a caller that actually reports the id ( + * {@code fleet_list}, the REST roster) sees one a lazy-resolving adapter (opencode) has since + * written, rather than the null frozen in at spawn time. Reserved for caller-driven reads, not + * timers: see {@link #roster()}'s javadoc for why the plain roster must stay non-resolving. + */ + public List rosterResolved() { + return registry.values().stream().map(this::resolveAgentSessionId).toList(); + } + + /** + * Resolve {@code session}'s {@code agentSessionId} if still unknown, re-polling the retained + * {@link PeerHandle} for this pane (fleetd #209). A no-op — returning {@code session} unchanged + * — once the id is already known, once no handle is retained for this pane (never spawned, or + * already released), or if the handle throws while answering. A resolved id is best-effort + * CAS-swapped into the registry via {@link #replace}; a lost race just means another caller + * already applied the same update, so the resolved value is returned either way. + */ + private MemberSession resolveAgentSessionId(MemberSession session) { + return resolveAgentSessionId(session, handles.get(session.paneId())); + } + + private MemberSession resolveAgentSessionId(MemberSession session, PeerHandle handle) { + if (session.agentSessionId() != null || handle == null) { + return session; + } + String resolved; + try { + resolved = handle.agentSessionId(); + } catch (RuntimeException e) { + log.debug("agentSessionId lookup failed for pane={} terminal={}: {}", + session.paneId(), session.terminalId(), e.toString()); + return session; + } + if (resolved == null) { + return session; + } + MemberSession updated = session.withAgentSessionId(resolved); + replace(session, updated); // best-effort; a lost CAS just means the resolved value stands anyway + return updated; + } + /** * 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. diff --git a/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java b/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java index 08f44b6..9fc0d37 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java @@ -931,4 +931,235 @@ class SessionManagerTest { assertTrue(e.getMessage().contains("SESSION_RESUME"), e.getMessage()); assertTrue(e.getMessage().contains("stub-profile"), e.getMessage()); } + + // ── fleetd #209: agentSessionId resolved lazily against the retained handle ───────────────── + // + // The opencode adapter cannot answer PeerHandle.agentSessionId() at spawn time — the on-disk + // session row is written only after the pane is live — so the id must be re-polled on a LATER + // call, against the SAME handle instance the launcher returned at spawn. SessionManager used + // to let that handle go out of scope at the end of the spawn method, so no caller ever re-asked + // it and fleet_list/fleet_status never saw the id. LazyIdHandle below reproduces exactly that + // shape: null on the first N calls (the spawn-time call included), a real id after. + + /** + * A {@link PeerHandle} whose {@link #agentSessionId()} answers {@code null} for its first + * {@code nullCalls} invocations, then either a fixed id or a configured throw on every call + * after that — the shape of the opencode bug (fleetd #209): the session row is not written + * until after the pane is live, so early polls come back empty and a later one finds it. + */ + private static final class LazyIdHandle implements PeerHandle { + private final String id; + private final String terminalId; + private final int nullCalls; + private final String resolvedId; + private final java.util.concurrent.atomic.AtomicInteger calls = + new java.util.concurrent.atomic.AtomicInteger(); + private volatile RuntimeException throwAfter; + + LazyIdHandle(String id, String terminalId, int nullCalls, String resolvedId) { + this.id = id; + this.terminalId = terminalId; + this.nullCalls = nullCalls; + this.resolvedId = resolvedId; + } + + /** After the null calls are exhausted, throw instead of answering the resolved id. */ + LazyIdHandle throwing(RuntimeException e) { + this.throwAfter = e; + return this; + } + + @Override + public String id() { + return id; + } + + @Override + public String terminalId() { + return terminalId; + } + + @Override + public String agentSessionId() { + int n = calls.incrementAndGet(); + if (n <= nullCalls) { + return null; + } + if (throwAfter != null) { + throw throwAfter; + } + return resolvedId; + } + + @Override + public CharterReceipt charterReceipt() { + return null; + } + + int callCount() { + return calls.get(); + } + } + + /** A minimal {@link PeerLauncher} that hands out pre-built {@link LazyIdHandle}s, one per spawn. */ + private static final class LazyIdLauncher implements PeerLauncher { + private final java.util.Deque queued = new java.util.ArrayDeque<>(); + + LazyIdLauncher queue(LazyIdHandle handle) { + queued.add(handle); + return this; + } + + @Override + public Set capabilities() { + return Set.of(Capability.WORKTREE, Capability.SESSION_RESUME); + } + + @Override + public Set capabilitiesFor(String profileName) { + return capabilities(); + } + + @Override + public PeerHandle spawn(SpawnRequest req) { + LazyIdHandle handle = queued.poll(); + if (handle == null) { + throw new IllegalStateException("no queued LazyIdHandle for this spawn"); + } + return handle; + } + + @Override + public Set profiles() { + return Set.of("lazy"); + } + + @Override + public String defaultProfile() { + return "lazy"; + } + + @Override + public String effectiveCwd(SpawnRequest req) { + return "/cwd"; + } + + @Override + public List parityOverlay(String profileName) { + return List.of(); + } + + @Override + public List list() { + return List.of(); + } + + @Override + public int reapOrphanWorkers() { + return 0; + } + + @Override + public void stop(String id) { + } + + @Override + public boolean clearContext(String id) { + return false; + } + } + + @Test + void plainRosterDoesNotResolveAgentSessionId() { + // fleetd #209 follow-up: roster() sits on the heartbeat/health-tick timers (and the metrics + // scrape), so it must never trigger the resolve lookup — for opencode that lookup opens an + // on-disk session database, and a member whose id never appears would pay that cost forever. + // rosterResolved() is the one to use when a caller actually reports the id. + LazyIdHandle handle = new LazyIdHandle("p0", "t0", 1, "oc-session-0"); + SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle)); + MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null); + assertEquals(1, handle.callCount(), "sanity: only the spawn-time call happened so far"); + + List roster = sessions.roster(); + + assertEquals(1, roster.size()); + assertEquals(acquired.paneId(), roster.getFirst().paneId()); + assertNull(roster.getFirst().agentSessionId(), "the plain roster must not resolve the id"); + assertEquals(1, handle.callCount(), + "roster() must never call agentSessionId() again — it sits on the heartbeat/health timers"); + } + + @Test + void rosterResolvedResolvesALateAgentSessionIdFromTheRetainedHandle() { + LazyIdHandle handle = new LazyIdHandle("p1", "t1", 1, "oc-session-1"); + SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle)); + + MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null); + assertNull(acquired.agentSessionId(), + "opencode has not written its session row yet at spawn time"); + + List roster = sessions.rosterResolved(); + assertEquals(1, roster.size()); + assertEquals("oc-session-1", roster.getFirst().agentSessionId(), + "fleet_list must see the id once the adapter can answer it"); + + Map view = SessionManager.rosterView(roster.getFirst(), null); + assertEquals("oc-session-1", view.get("agentSessionId"), + "rosterView renders whatever rosterResolved() resolved"); + } + + @Test + void getResolvesALateAgentSessionIdFromTheRetainedHandle() { + LazyIdHandle handle = new LazyIdHandle("p2", "t2", 1, "oc-session-2"); + SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle)); + MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null); + + MemberSession resolved = sessions.get(acquired.paneId()).orElseThrow(); + + assertEquals("oc-session-2", resolved.agentSessionId(), + "fleet_status (single-session lookup) must also see the late-resolved id"); + } + + @Test + void releaseCarriesALateResolvedAgentSessionIdIntoTheReleaseDetail() { + LazyIdHandle handle = new LazyIdHandle("p3", "t3", 1, "oc-session-3"); + SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle)); + MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null); + java.util.List released = new java.util.concurrent.CopyOnWriteArrayList<>(); + sessions.onRelease(detail -> released.add(detail.agentSessionId())); + + sessions.release(acquired.paneId()); + + assertEquals(java.util.List.of("oc-session-3"), released, + "a released member's detail carries the id it has since resolved, not the null " + + "frozen in at spawn time"); + } + + @Test + void aThrowingHandleDoesNotBreakRosterResolved() { + LazyIdHandle handle = new LazyIdHandle("p4", "t4", 1, "oc-session-4") + .throwing(new RuntimeException("sqlite locked")); + SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle)); + sessions.acquire("lazy", "/cwd", "/caller", null); + + List roster = assertDoesNotThrow(sessions::rosterResolved, + "a handle that throws resolving its id must not break the roster read"); + + assertEquals(1, roster.size()); + assertNull(roster.getFirst().agentSessionId(), "the id stays unresolved when the lookup throws"); + } + + @Test + void aResolvedAgentSessionIdIsNotLookedUpAgain() { + LazyIdHandle handle = new LazyIdHandle("p5", "t5", 1, "oc-session-5"); + SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle)); + sessions.acquire("lazy", "/cwd", "/caller", null); + assertEquals(1, handle.callCount(), "sanity: only the spawn-time call happened so far"); + + sessions.rosterResolved(); + assertEquals(2, handle.callCount(), "the first rosterResolved() read resolves the id"); + + sessions.rosterResolved(); + assertEquals(2, handle.callCount(), "once resolved, the id must not be looked up again"); + } }