#209: resolve agentSessionId lazily via the retained PeerHandle
SessionManager called handle.agentSessionId() once at spawn and froze it in the
immutable MemberSession. For opencode that value is always null - the session
row does not exist yet when the pane is created - so fleet_list never reported
an agentSessionId and fleet_spawn{resumeSessionId} was unusable for that
backend. The handle's javadoc said 'the caller re-calls later'; no caller did,
and SessionManager did not even retain the handle.
Retain the PeerHandle per pane and re-resolve while the stored id is still
null, sticky once found, CAS-swapped into the registry.
roster() stays non-resolving - it is the roster supplier for LeadHeartbeatLoop
and FleetHealthMonitor, and resolving there would open opencode's 841MB SQLite
database on every tick, for every unresolved member, forever. rosterResolved()
carries the resolve and is used only by fleet_list and the REST roster, the two
surfaces that report the id. get(paneId) and release() resolve too, both
caller-driven.
A test pins the split: the plain roster() must never call agentSessionId()
again.
This commit is contained in:
@@ -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<MemberSession> 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<MemberSession> roster = sessions.rosterResolved();
|
||||
List<Map<String, Object>> out = roster.stream()
|
||||
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
|
||||
.toList();
|
||||
|
||||
@@ -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<Map<String, Object>> 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<Map<String, Object>> out = sessions.rosterResolved().stream()
|
||||
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
||||
.toList();
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -46,6 +46,15 @@ public final class SessionManager implements TurnListener {
|
||||
private final PeerLauncher launcher;
|
||||
private final Worktrees worktrees;
|
||||
private final ConcurrentHashMap<String /*paneId*/, MemberSession> 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<String /*paneId*/, PeerHandle> 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<MemberSession> 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
|
||||
* <strong>not</strong> 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<MemberSession> 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<MemberSession> 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.
|
||||
|
||||
@@ -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<LazyIdHandle> queued = new java.util.ArrayDeque<>();
|
||||
|
||||
LazyIdLauncher queue(LazyIdHandle handle) {
|
||||
queued.add(handle);
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return Set.of(Capability.WORKTREE, Capability.SESSION_RESUME);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> 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<String> profiles() {
|
||||
return Set.of("lazy");
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return "lazy";
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return "/cwd";
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> 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<MemberSession> 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<MemberSession> 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<String, Object> 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<String> 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<MemberSession> 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");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user