diff --git a/hub/src/sync/sessionCache.ts b/hub/src/sync/sessionCache.ts index ea4d3d57..902bda9a 100644 --- a/hub/src/sync/sessionCache.ts +++ b/hub/src/sync/sessionCache.ts @@ -10,6 +10,7 @@ export class SessionCache { private readonly sessions: Map = new Map() private readonly lastBroadcastAtBySessionId: Map = new Map() private readonly todoBackfillAttemptedSessionIds: Set = new Set() + private readonly deduplicateInProgress: Set = new Set() constructor( private readonly store: Store, @@ -280,16 +281,20 @@ export class SessionCache { this.publisher.emit({ type: 'session-updated', sessionId: session.id, data: { active: false, thinking: false, backgroundTaskCount: 0 } }) } - expireInactive(now: number = Date.now()): void { + expireInactive(now: number = Date.now()): string[] { const sessionTimeoutMs = 30_000 + const expired: string[] = [] for (const session of this.sessions.values()) { if (!session.active) continue if (now - session.activeAt <= sessionTimeoutMs) continue session.active = false session.thinking = false + expired.push(session.id) this.publisher.emit({ type: 'session-updated', sessionId: session.id, data: { active: false } }) } + + return expired } applySessionConfig( @@ -470,6 +475,27 @@ export class SessionCache { ) } + // Merge agentState: union requests/completedRequests from both sessions so pending + // approvals on the duplicate are not lost. Only inactive duplicates reach this point + // (active ones are skipped by deduplicateByAgentSessionId). + // Read the latest target state right before writing to avoid overwriting live updates. + if (oldStored.agentState !== null) { + for (let attempt = 0; attempt < 2; attempt += 1) { + const latest = this.store.sessions.getSessionByNamespace(newSessionId, namespace) + if (!latest) break + const mergedAgentState = this.mergeAgentState(oldStored.agentState, latest.agentState) + if (mergedAgentState === null || mergedAgentState === latest.agentState) break + const result = this.store.sessions.updateSessionAgentState( + newSessionId, + mergedAgentState, + latest.agentStateVersion, + namespace + ) + if (result.result !== 'version-mismatch') break + // version-mismatch: retry with fresh snapshot + } + } + if (oldStored.teamState !== null && oldStored.teamStateUpdatedAt !== null) { this.store.sessions.setSessionTeamState( newSessionId, @@ -537,4 +563,86 @@ export class SessionCache { return changed ? merged : newMetadata } + + private mergeAgentState(oldState: unknown | null, newState: unknown | null): unknown | null { + if (oldState === null) return newState + if (newState === null) return oldState + + const oldObj = oldState as Record + const newObj = newState as Record + + const completedRequests = { + ...((oldObj.completedRequests as Record | undefined) ?? {}), + ...((newObj.completedRequests as Record | undefined) ?? {}) + } + // Filter out requests that are already completed to avoid resurrecting them as pending + const completedIds = new Set(Object.keys(completedRequests)) + const requests = Object.fromEntries( + Object.entries({ + ...((oldObj.requests as Record | undefined) ?? {}), + ...((newObj.requests as Record | undefined) ?? {}) + }).filter(([id]) => !completedIds.has(id)) + ) + + return { ...oldObj, ...newObj, requests, completedRequests } + } + + private extractAgentSessionId( + metadata: NonNullable + ): { field: 'codexSessionId' | 'claudeSessionId' | 'geminiSessionId' | 'opencodeSessionId' | 'cursorSessionId'; value: string } | null { + if (metadata.codexSessionId) return { field: 'codexSessionId', value: metadata.codexSessionId } + if (metadata.claudeSessionId) return { field: 'claudeSessionId', value: metadata.claudeSessionId } + if (metadata.geminiSessionId) return { field: 'geminiSessionId', value: metadata.geminiSessionId } + if (metadata.opencodeSessionId) return { field: 'opencodeSessionId', value: metadata.opencodeSessionId } + if (metadata.cursorSessionId) return { field: 'cursorSessionId', value: metadata.cursorSessionId } + return null + } + + async deduplicateByAgentSessionId(sessionId: string): Promise { + const session = this.sessions.get(sessionId) + if (!session?.metadata) return + + const agentId = this.extractAgentSessionId(session.metadata) + if (!agentId) return + + // Guard: skip if another dedup for this agent ID is already in progress. + // A skipped trigger is acceptable — the web-side display dedup hides any remaining duplicates. + if (this.deduplicateInProgress.has(agentId.value)) return + this.deduplicateInProgress.add(agentId.value) + + try { + const candidates: { id: string; session: Session }[] = [{ id: sessionId, session }] + for (const [existingId, existing] of this.sessions) { + if (existingId === sessionId) continue + if (existing.namespace !== session.namespace) continue + if (!existing.metadata) continue + if (existing.metadata[agentId.field] !== agentId.value) continue + // Only merge inactive duplicates. Active ones still have a live CLI socket + // whose keepalive/messages would fail if we deleted their session record. + // The web-side display dedup hides active duplicates from the UI. + if (existing.active) continue + candidates.push({ id: existingId, session: existing }) + } + + if (candidates.length <= 1) return + + // Keep the most recent session as the merge target so newer state survives. + candidates.sort((a, b) => + (b.session.activeAt - a.session.activeAt) || (b.session.updatedAt - a.session.updatedAt) + ) + const targetId = candidates[0].id + const targetNamespace = candidates[0].session.namespace + + for (const { id } of candidates.slice(1)) { + if (id === targetId) continue + try { + await this.mergeSessions(id, targetId, targetNamespace) + } catch { + // best-effort: duplicate remains if merge fails + } + } + } finally { + this.deduplicateInProgress.delete(agentId.value) + } + } } diff --git a/hub/src/sync/sessionModel.test.ts b/hub/src/sync/sessionModel.test.ts index a923118b..9b7642d9 100644 --- a/hub/src/sync/sessionModel.test.ts +++ b/hub/src/sync/sessionModel.test.ts @@ -440,4 +440,262 @@ describe('session model', () => { engine.stop() } }) + + describe('session dedup by agent session ID', () => { + it('merges duplicate when codexSessionId collides', async () => { + const store = new Store(':memory:') + const events: SyncEvent[] = [] + const cache = new SessionCache(store, createPublisher(events)) + + const s1 = cache.getOrCreateSession( + 'tag-1', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + + // Add a message to s1 + store.messages.addMessage(s1.id, { type: 'text', text: 'hello from s1' }, 'local-1') + + const s2 = cache.getOrCreateSession( + 'tag-2', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + + expect(s1.id).not.toBe(s2.id) + + await cache.deduplicateByAgentSessionId(s2.id) + + expect(cache.getSession(s1.id)).toBeUndefined() + expect(cache.getSession(s2.id)).toBeDefined() + + const messages = store.messages.getMessages(s2.id, 100) + expect(messages.length).toBeGreaterThanOrEqual(1) + }) + + it('preserves sessions with different agent session IDs', async () => { + const store = new Store(':memory:') + const events: SyncEvent[] = [] + const cache = new SessionCache(store, createPublisher(events)) + + const s1 = cache.getOrCreateSession( + 'tag-1', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + const s2 = cache.getOrCreateSession( + 'tag-2', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-Y' }, + null, + 'default' + ) + + await cache.deduplicateByAgentSessionId(s2.id) + + expect(cache.getSession(s1.id)).toBeDefined() + expect(cache.getSession(s2.id)).toBeDefined() + }) + + it('does not merge across namespaces', async () => { + const store = new Store(':memory:') + const events: SyncEvent[] = [] + const cache = new SessionCache(store, createPublisher(events)) + + const s1 = cache.getOrCreateSession( + 'tag-1', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'ns1' + ) + const s2 = cache.getOrCreateSession( + 'tag-2', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'ns2' + ) + + await cache.deduplicateByAgentSessionId(s2.id) + + expect(cache.getSession(s1.id)).toBeDefined() + expect(cache.getSession(s2.id)).toBeDefined() + }) + + it('no-op when session has no agent session ID', async () => { + const store = new Store(':memory:') + const events: SyncEvent[] = [] + const cache = new SessionCache(store, createPublisher(events)) + + const s1 = cache.getOrCreateSession( + 'tag-1', + { path: '/tmp/project', host: 'localhost', flavor: 'codex' }, + null, + 'default' + ) + + await cache.deduplicateByAgentSessionId(s1.id) + + expect(cache.getSession(s1.id)).toBeDefined() + }) + + it('does not merge active duplicates', async () => { + const store = new Store(':memory:') + const events: SyncEvent[] = [] + const cache = new SessionCache(store, createPublisher(events)) + + const s1 = cache.getOrCreateSession( + 'tag-1', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + + // Mark s1 as active (simulating a live CLI connection) + cache.handleSessionAlive({ sid: s1.id, time: Date.now(), thinking: false }) + + const s2 = cache.getOrCreateSession( + 'tag-2', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + + await cache.deduplicateByAgentSessionId(s2.id) + + // s1 is active, so it should NOT be merged/deleted + expect(cache.getSession(s1.id)).toBeDefined() + expect(cache.getSession(s2.id)).toBeDefined() + }) + + it('merges duplicate after it becomes inactive via session-end', async () => { + const store = new Store(':memory:') + const engine = new SyncEngine( + store, + {} as never, + new RpcRegistry(), + { broadcast() {} } as never + ) + + try { + const s1 = engine.getOrCreateSession( + 'tag-1', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + const s2 = engine.getOrCreateSession( + 'tag-2', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + + // Mark s1 as active + engine.handleSessionAlive({ sid: s1.id, time: Date.now() }) + + // s1 is active, dedup from s2 should skip it + const events: SyncEvent[] = [] + const cache = (engine as any).sessionCache as SessionCache + await cache.deduplicateByAgentSessionId(s2.id) + expect(cache.getSession(s1.id)).toBeDefined() + expect(cache.getSession(s2.id)).toBeDefined() + + // Now s1 ends — handleSessionEnd should trigger dedup retry + engine.handleSessionEnd({ sid: s1.id, time: Date.now() }) + + // Give the fire-and-forget dedup a tick to complete + await new Promise((r) => setTimeout(r, 50)) + + // One of them should be merged away + const s1Exists = cache.getSession(s1.id) + const s2Exists = cache.getSession(s2.id) + expect(!s1Exists || !s2Exists).toBe(true) + } finally { + engine.stop() + } + }) + + it('merges duplicate after inactivity timeout expires it', async () => { + const store = new Store(':memory:') + const events: SyncEvent[] = [] + const cache = new SessionCache(store, createPublisher(events)) + + const s1 = cache.getOrCreateSession( + 'tag-1', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + const s2 = cache.getOrCreateSession( + 'tag-2', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + null, + 'default' + ) + + // Mark s1 as active now + cache.handleSessionAlive({ sid: s1.id, time: Date.now() }) + + // s1 is active — dedup skips it + await cache.deduplicateByAgentSessionId(s2.id) + expect(cache.getSession(s1.id)).toBeDefined() + + // Simulate time passing beyond the 30s timeout + const expired = cache.expireInactive(Date.now() + 60_000) + expect(expired).toContain(s1.id) + + // Now s1 is inactive — dedup should merge it + await cache.deduplicateByAgentSessionId(s2.id) + expect(cache.getSession(s1.id)).toBeUndefined() + expect(cache.getSession(s2.id)).toBeDefined() + }) + + it('deep-merges agentState and filters completed requests', async () => { + const store = new Store(':memory:') + const events: SyncEvent[] = [] + const cache = new SessionCache(store, createPublisher(events)) + + const s1 = cache.getOrCreateSession( + 'tag-1', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + { + requests: { + 'req-1': { tool: 'Bash', arguments: {} }, + 'req-2': { tool: 'Bash', arguments: {} } + }, + completedRequests: {} + }, + 'default' + ) + const s2 = cache.getOrCreateSession( + 'tag-2', + { path: '/tmp/project', host: 'localhost', flavor: 'codex', codexSessionId: 'thread-X' }, + { + requests: { + 'req-3': { tool: 'Bash', arguments: {} } + }, + completedRequests: { + 'req-1': { tool: 'Bash', arguments: {}, status: 'approved' } + } + }, + 'default' + ) + + await cache.deduplicateByAgentSessionId(s2.id) + + const session = cache.getSession(s2.id) + expect(session).toBeDefined() + const state = session!.agentState! + + // req-1 was completed in s2 — should NOT appear in requests + expect(state.requests?.['req-1']).toBeUndefined() + // req-2 and req-3 are still pending + expect(state.requests?.['req-2']).toBeDefined() + expect(state.requests?.['req-3']).toBeDefined() + // completedRequests has req-1 + expect(state.completedRequests?.['req-1']).toBeDefined() + }) + }) }) diff --git a/hub/src/sync/syncEngine.ts b/hub/src/sync/syncEngine.ts index 11275952..9708ae2b 100644 --- a/hub/src/sync/syncEngine.ts +++ b/hub/src/sync/syncEngine.ts @@ -163,7 +163,16 @@ export class SyncEngine { handleRealtimeEvent(event: SyncEvent): void { if (event.type === 'session-updated' && event.sessionId) { + // Snapshot agent session IDs before refresh — safe because JS is single-threaded + // and refreshSession replaces the Map entry with a new object. + const before = this.sessionCache.getSession(event.sessionId) this.sessionCache.refreshSession(event.sessionId) + const after = this.sessionCache.getSession(event.sessionId) + if (after?.metadata && !this.hasSameAgentSessionIds(before?.metadata ?? null, after.metadata)) { + void this.sessionCache.deduplicateByAgentSessionId(event.sessionId).catch(() => { + // best-effort: dedup failure is harmless, web-side safety net hides remaining duplicates + }) + } return } @@ -197,6 +206,9 @@ export class SyncEngine { handleSessionEnd(payload: { sid: string; time: number }): void { this.sessionCache.handleSessionEnd(payload) + // Retry dedup now that this session is inactive — a prior dedup may have + // skipped it because it was still active at the time. + this.triggerDedupIfNeeded(payload.sid) } handleBackgroundTaskDelta(sessionId: string, delta: { started: number; completed: number }): void { @@ -208,7 +220,16 @@ export class SyncEngine { } private expireInactive(): void { - this.sessionCache.expireInactive() + const expired = this.sessionCache.expireInactive() + // Sort by most recent first so dedup keeps the newest session when multiple + // duplicates for the same agent thread expire in the same sweep. + const sorted = expired + .map((id) => this.sessionCache.getSession(id)) + .filter((s): s is NonNullable => s != null) + .sort((a, b) => (b.activeAt - a.activeAt) || (b.updatedAt - a.updatedAt)) + for (const session of sorted) { + this.triggerDedupIfNeeded(session.id) + } this.machineCache.expireInactive() } @@ -430,17 +451,43 @@ export class SyncEngine { } if (spawnResult.sessionId !== access.sessionId) { - try { - await this.sessionCache.mergeSessions(access.sessionId, spawnResult.sessionId, namespace) - } catch (error) { - const message = error instanceof Error ? error.message : 'Failed to merge resumed session' - return { type: 'error', message, code: 'resume_failed' } + // The old session may have already been merged by the automatic dedup path + // (triggered when the spawned CLI sets its agent session ID in metadata). + // Only attempt the explicit merge if the old session still exists. + const oldSession = this.sessionCache.getSessionByNamespace(access.sessionId, namespace) + if (oldSession) { + try { + await this.sessionCache.mergeSessions(access.sessionId, spawnResult.sessionId, namespace) + } catch (error) { + const message = error instanceof Error ? error.message : 'Failed to merge resumed session' + return { type: 'error', message, code: 'resume_failed' } + } } } return { type: 'success', sessionId: spawnResult.sessionId } } + private hasSameAgentSessionIds( + prev: Session['metadata'] | null, + next: NonNullable + ): boolean { + return (prev?.codexSessionId ?? null) === (next.codexSessionId ?? null) + && (prev?.claudeSessionId ?? null) === (next.claudeSessionId ?? null) + && (prev?.geminiSessionId ?? null) === (next.geminiSessionId ?? null) + && (prev?.opencodeSessionId ?? null) === (next.opencodeSessionId ?? null) + && (prev?.cursorSessionId ?? null) === (next.cursorSessionId ?? null) + } + + private triggerDedupIfNeeded(sessionId: string): void { + const session = this.sessionCache.getSession(sessionId) + if (session?.metadata) { + void this.sessionCache.deduplicateByAgentSessionId(sessionId).catch(() => { + // best-effort: web-side safety net hides remaining duplicates + }) + } + } + async waitForSessionActive(sessionId: string, timeoutMs: number = 15_000): Promise { const start = Date.now() while (Date.now() - start < timeoutMs) { diff --git a/shared/src/sessionSummary.ts b/shared/src/sessionSummary.ts index e717a57d..86298208 100644 --- a/shared/src/sessionSummary.ts +++ b/shared/src/sessionSummary.ts @@ -7,6 +7,7 @@ export type SessionSummaryMetadata = { summary?: { text: string } flavor?: string | null worktree?: WorktreeMetadata + agentSessionId?: string } export type SessionSummary = { @@ -31,7 +32,13 @@ export function toSessionSummary(session: Session): SessionSummary { machineId: session.metadata.machineId ?? undefined, summary: session.metadata.summary ? { text: session.metadata.summary.text } : undefined, flavor: session.metadata.flavor ?? null, - worktree: session.metadata.worktree + worktree: session.metadata.worktree, + agentSessionId: session.metadata.codexSessionId + ?? session.metadata.claudeSessionId + ?? session.metadata.geminiSessionId + ?? session.metadata.opencodeSessionId + ?? session.metadata.cursorSessionId + ?? undefined } : null const todoProgress = session.todos?.length ? { diff --git a/web/src/components/SessionList.test.ts b/web/src/components/SessionList.test.ts new file mode 100644 index 00000000..b830e801 --- /dev/null +++ b/web/src/components/SessionList.test.ts @@ -0,0 +1,82 @@ +import { describe, expect, it } from 'vitest' +import type { SessionSummary } from '@/types/api' +import { deduplicateSessionsByAgentId } from './SessionList' + +function makeSession(overrides: Partial & { id: string }): SessionSummary { + return { + active: false, + thinking: false, + activeAt: 0, + updatedAt: 0, + metadata: null, + todoProgress: null, + pendingRequestsCount: 0, + model: null, + effort: null, + ...overrides + } +} + +describe('deduplicateSessionsByAgentId', () => { + it('deduplicates sessions with the same agentSessionId', () => { + const sessions = [ + makeSession({ id: 'a', metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 100 }), + makeSession({ id: 'b', metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 200 }) + ] + const result = deduplicateSessionsByAgentId(sessions) + expect(result).toHaveLength(1) + expect(result[0].id).toBe('b') // more recent wins + }) + + it('keeps active session over inactive duplicate', () => { + const sessions = [ + makeSession({ id: 'a', active: true, metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 100 }), + makeSession({ id: 'b', metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 200 }) + ] + const result = deduplicateSessionsByAgentId(sessions) + expect(result).toHaveLength(1) + expect(result[0].id).toBe('a') // active wins despite older updatedAt + }) + + it('prefers selected session among inactive duplicates', () => { + const sessions = [ + makeSession({ id: 'a', metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 100 }), + makeSession({ id: 'b', metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 200 }) + ] + const result = deduplicateSessionsByAgentId(sessions, 'a') + expect(result).toHaveLength(1) + expect(result[0].id).toBe('a') // selected wins despite older updatedAt + }) + + it('active always wins over selected inactive', () => { + const sessions = [ + makeSession({ id: 'a', metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 200 }), + makeSession({ id: 'b', active: true, metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 100 }) + ] + const result = deduplicateSessionsByAgentId(sessions, 'a') + expect(result).toHaveLength(1) + expect(result[0].id).toBe('b') // active wins over selected + }) + + it('passes through sessions without agentSessionId', () => { + const sessions = [ + makeSession({ id: 'a', metadata: { path: '/p' } }), + makeSession({ id: 'b', metadata: { path: '/p', agentSessionId: 'thread-1' } }), + makeSession({ id: 'c', metadata: null }) + ] + const result = deduplicateSessionsByAgentId(sessions) + expect(result).toHaveLength(3) + }) + + it('deduplicates independently across different agentSessionIds', () => { + const sessions = [ + makeSession({ id: 'a', metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 100 }), + makeSession({ id: 'b', metadata: { path: '/p', agentSessionId: 'thread-1' }, updatedAt: 200 }), + makeSession({ id: 'c', metadata: { path: '/p', agentSessionId: 'thread-2' }, updatedAt: 100 }), + makeSession({ id: 'd', metadata: { path: '/p', agentSessionId: 'thread-2' }, updatedAt: 200 }) + ] + const result = deduplicateSessionsByAgentId(sessions) + expect(result).toHaveLength(2) + expect(result.map(s => s.id).sort()).toEqual(['b', 'd']) + }) +}) diff --git a/web/src/components/SessionList.tsx b/web/src/components/SessionList.tsx index 9fef3c97..115cc266 100644 --- a/web/src/components/SessionList.tsx +++ b/web/src/components/SessionList.tsx @@ -39,6 +39,39 @@ function getGroupDisplayName(directory: string): string { export const UNKNOWN_MACHINE_ID = '__unknown__' +export function deduplicateSessionsByAgentId(sessions: SessionSummary[], selectedSessionId?: string | null): SessionSummary[] { + const byAgentId = new Map() + const result: SessionSummary[] = [] + + for (const session of sessions) { + const agentId = session.metadata?.agentSessionId + if (!agentId) { + result.push(session) + continue + } + const group = byAgentId.get(agentId) + if (group) { + group.push(session) + } else { + byAgentId.set(agentId, [session]) + } + } + + for (const group of byAgentId.values()) { + group.sort((a, b) => { + // Active session always wins — it's the live connection + if (a.active !== b.active) return a.active ? -1 : 1 + // Among inactive duplicates, keep the selected one visible + if (a.id === selectedSessionId) return -1 + if (b.id === selectedSessionId) return 1 + return b.updatedAt - a.updatedAt + }) + result.push(group[0]) + } + + return result +} + function groupSessionsByDirectory(sessions: SessionSummary[]): SessionGroup[] { const groups = new Map() @@ -453,8 +486,8 @@ export function SessionList(props: { const { t } = useTranslation() const { renderHeader = true, api, selectedSessionId, machineLabelsById = {} } = props const groups = useMemo( - () => groupSessionsByDirectory(props.sessions), - [props.sessions] + () => groupSessionsByDirectory(deduplicateSessionsByAgentId(props.sessions, selectedSessionId)), + [props.sessions, selectedSessionId] ) const [collapseOverrides, setCollapseOverrides] = useState>( () => new Map()