import { DisposableStore } from '#/_base/di/lifecycle'; import { Emitter, type Event, type IWaitUntil } from '#/_base/event'; import { ScopeActivation, registerScopedService, type ISessionScopeHandle } from '#/_base/di/scope'; import { LifecycleScope } from '#/app/scopes'; import { Error2, ErrorCodes } from '#/errors'; import { ISessionIndex, type SessionSummary } from '#/app/sessionIndex/sessionIndex'; import type { SessionMeta } from '#/session/sessionMetadata/sessionMetadata'; import { type CreateChildSessionOptions, type ForkSessionOptions, type ResumeSessionOptions, type SessionArchivedEvent, type SessionClosedEvent, type SessionCreatedEvent, type SessionForkedEvent, type SessionWillCloseEvent, type SessionWillCreateEvent, } from '#/workspace/sessionLifecycle/sessionLifecycle'; import type { SessionLifecycleService } from '#/workspace/sessionLifecycle/sessionLifecycleService'; import { IWorkspaceInstanceManager } from '#/workspace/workspaceInstance/workspaceInstanceManager'; import { ISessionManager, type CreateManagedSessionOptions, type UnguardedSessionLifecycle, } from './sessionManager'; interface SessionControllerEntry { readonly generation: string; readonly controller: SessionLifecycleService; readonly subscriptions: DisposableStore; sessionCount: number; } export class SessionManager implements ISessionManager { declare readonly _serviceBrand: undefined; private readonly sessions = new Map(); private readonly owners = new Map(); private readonly pendingResumes = new Map>(); private readonly resumeFailures = new Map(); private readonly lifecycleChains = new Map>(); private readonly controllers = new Map(); private readonly controllerEntries = new Set(); private readonly willCreateEmitter = new Emitter(); readonly onWillCreateSession: Event = this.willCreateEmitter.event; private readonly didCreateEmitter = new Emitter(); readonly onDidCreateSession = this.didCreateEmitter.event; private readonly willCloseEmitter = new Emitter(); readonly onWillCloseSession = this.willCloseEmitter.event; private readonly didCloseEmitter = new Emitter(); readonly onDidCloseSession = this.didCloseEmitter.event; private readonly willDeleteEmitter = new Emitter<{ readonly sessionId: string } & IWaitUntil>(); readonly onWillDeleteSession = this.willDeleteEmitter.event; private readonly didArchiveEmitter = new Emitter(); readonly onDidArchiveSession = this.didArchiveEmitter.event; private readonly didForkEmitter = new Emitter(); readonly onDidForkSession = this.didForkEmitter.event; constructor( @IWorkspaceInstanceManager private readonly workspaces: IWorkspaceInstanceManager, @ISessionIndex private readonly index: ISessionIndex, ) {} async create(options: CreateManagedSessionOptions): Promise { const workspace = await this.workspaces.getOrCreate( options.workspaceId === undefined ? { root: options.workDir } : { workspaceId: options.workspaceId, root: options.workDir }, ); const create = () => this.controllerForWorkspace(workspace.id).create(options); if (options.sessionId === undefined) return create(); return this.serializeLifecycle(options.sessionId, create); } async resume(sessionId: string, options?: ResumeSessionOptions): Promise { const inflight = this.pendingResumes.get(sessionId); if (inflight !== undefined) return inflight; this.resumeFailures.delete(sessionId); const promise = this.serializeLifecycle(sessionId, async () => (await this.controllerForSession(sessionId))?.resume(sessionId, options), ).finally(() => this.pendingResumes.delete(sessionId)); this.pendingResumes.set(sessionId, promise); void promise.catch((error: unknown) => { this.resumeFailures.set(sessionId, error instanceof Error ? error : new Error('session resume failed')); }); return promise; } get(sessionId: string): ISessionScopeHandle | undefined { return this.sessions.get(sessionId); } status(sessionId: string): Promise { return this.index.get(sessionId); } async whenResumeSettled(sessionId: string): Promise { await this.pendingResumes.get(sessionId); const failure = this.resumeFailures.get(sessionId); if (failure !== undefined) throw failure; await this.owners.get(sessionId)?.whenResumeSettled(sessionId); } private serializeLifecycle(sessionId: string, work: () => Promise): Promise { const prev = this.lifecycleChains.get(sessionId) ?? Promise.resolve(); const run = prev.then(work, work); const next = run.then( () => undefined, () => undefined, ); this.lifecycleChains.set(sessionId, next); void next.finally(() => { if (this.lifecycleChains.get(sessionId) === next) this.lifecycleChains.delete(sessionId); }); return run; } private serializeLifecycleForKeys(keys: readonly string[], work: () => Promise): Promise { const [first, ...rest] = keys; if (first === undefined) return work(); return this.serializeLifecycle(first, () => this.serializeLifecycleForKeys(rest, work)); } private lifecycleKeys(...ids: (string | undefined)[]): string[] { return [...new Set(ids.filter((id): id is string => id !== undefined))].sort(); } withLifecycleSerialization( sessionId: string, work: (unguarded: UnguardedSessionLifecycle) => Promise, ): Promise { return this.serializeLifecycle(sessionId, () => work({ archive: () => this.archiveInner(sessionId), restore: () => this.restoreInner(sessionId), }), ); } list(): readonly ISessionScopeHandle[] { return [...this.sessions.values()]; } async close(sessionId: string): Promise { await this.serializeLifecycle(sessionId, async () => this.owners.get(sessionId)?.close(sessionId)); } private async archiveInner(sessionId: string): Promise { await (await this.controllerForSession(sessionId))?.archive(sessionId); } async archive(sessionId: string): Promise { await this.serializeLifecycle(sessionId, () => this.archiveInner(sessionId)); } private async restoreInner( sessionId: string, options?: ResumeSessionOptions, ): Promise { return (await this.controllerForSession(sessionId))?.restore(sessionId, options); } async restore(sessionId: string, options?: ResumeSessionOptions): Promise { return this.serializeLifecycle(sessionId, () => this.restoreInner(sessionId, options)); } async delete(sessionId: string): Promise { await this.serializeLifecycle(sessionId, async () => { const controller = await this.controllerForSession(sessionId); if (controller === undefined) { throw new Error2(ErrorCodes.SESSION_NOT_FOUND, `session ${sessionId} does not exist`); } await controller.close(sessionId); const cleanups: Promise[] = []; this.willDeleteEmitter.fire({ sessionId, signal: new AbortController().signal, waitUntil: (cleanup) => { if (Object.isFrozen(cleanups)) throw new Error('waitUntil must be called synchronously'); cleanups.push(cleanup); }, }); void Object.freeze(cleanups); const settled = await Promise.allSettled(cleanups); const failed = settled.find((result) => result.status === 'rejected'); if (failed?.status === 'rejected') throw failed.reason; await controller.delete(sessionId); }); } async fork(options: ForkSessionOptions): Promise { return this.serializeLifecycleForKeys( this.lifecycleKeys(options.sourceSessionId, options.newSessionId), async () => { const controller = await this.controllerForSession(options.sourceSessionId); if (controller === undefined) { throw new Error2( ErrorCodes.SESSION_NOT_FOUND, `session ${options.sourceSessionId} does not exist`, ); } return controller.fork(options); }, ); } async createChild(options: CreateChildSessionOptions): Promise { return this.serializeLifecycleForKeys( this.lifecycleKeys(options.sourceSessionId, options.newSessionId), async () => { const controller = await this.controllerForSession(options.sourceSessionId); if (controller === undefined) { throw new Error2( ErrorCodes.SESSION_NOT_FOUND, `session ${options.sourceSessionId} does not exist`, ); } return controller.createChild(options); }, ); } dispose(): void { for (const { controller, subscriptions } of [...this.controllerEntries].reverse()) { subscriptions.dispose(); controller.dispose(); } this.controllerEntries.clear(); this.controllers.clear(); this.sessions.clear(); this.owners.clear(); this.willCreateEmitter.dispose(); this.didCreateEmitter.dispose(); this.willCloseEmitter.dispose(); this.didCloseEmitter.dispose(); this.willDeleteEmitter.dispose(); this.didArchiveEmitter.dispose(); this.didForkEmitter.dispose(); } private controllerForWorkspace(workspaceId: string): SessionLifecycleService { const workspace = this.workspaces.get(workspaceId); if (workspace === undefined) throw new Error(`workspace ${workspaceId} is not materialized`); const generation = workspace.program.sessionControllerGeneration; const existing = this.controllers.get(workspaceId); if (existing?.generation === generation) return existing.controller; const controller = workspace.program.createSessionController(); const subscriptions = new DisposableStore(); const entry: SessionControllerEntry = { generation, controller, subscriptions, sessionCount: 0 }; subscriptions.add(controller.onWillCreateSession((event) => this.willCreateEmitter.fire(event))); subscriptions.add(controller.onDidCreateSession((event) => { entry.sessionCount += 1; this.sessions.set(event.sessionId, event.handle); this.owners.set(event.sessionId, controller); this.didCreateEmitter.fire(event); })); subscriptions.add(controller.onWillCloseSession((event) => this.willCloseEmitter.fire(event))); subscriptions.add(controller.onDidCloseSession((event) => { entry.sessionCount -= 1; this.sessions.delete(event.sessionId); this.owners.delete(event.sessionId); this.didCloseEmitter.fire(event); this.retireEntryIfIdle(workspaceId, entry); })); subscriptions.add(controller.onDidArchiveSession((event) => { entry.sessionCount -= 1; this.sessions.delete(event.sessionId); this.owners.delete(event.sessionId); this.didArchiveEmitter.fire(event); this.retireEntryIfIdle(workspaceId, entry); })); subscriptions.add(controller.onDidForkSession((event) => this.didForkEmitter.fire(event))); this.controllerEntries.add(entry); this.controllers.set(workspaceId, entry); if (existing !== undefined) this.retireEntryIfIdle(workspaceId, existing); return controller; } private retireEntryIfIdle(workspaceId: string, entry: SessionControllerEntry): void { if (entry.sessionCount !== 0 || !this.controllerEntries.has(entry)) return; this.controllerEntries.delete(entry); if (this.controllers.get(workspaceId) === entry) this.controllers.delete(workspaceId); entry.subscriptions.dispose(); entry.controller.dispose(); } private async controllerForSession(sessionId: string): Promise { const live = this.owners.get(sessionId); if (live !== undefined) return live; const summary = await this.index.get(sessionId); if (summary === undefined) return undefined; const workspace = await this.workspaces.getOrCreate({ workspaceId: summary.workspaceId, root: summary.cwd }); return this.controllerForWorkspace(workspace.id); } } registerScopedService(LifecycleScope.App, ISessionManager, SessionManager, ScopeActivation.OnScopeCreated, 'sessionManager');