Download packages/agent-core-v2/src/app/sessionManager/sessionManagerService.ts from SaylorTwift/kimi-code: direct link, hf CLI and curl.
- Browser
- Download file 12.7 kB
-
https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/agent-core-v2/src/app/sessionManager/sessionManagerService.ts
- Command line
-
hf download hf://SaylorTwift/kimi-code/packages/agent-core-v2/src/app/sessionManager/sessionManagerService.ts
-
curl -L -o sessionManagerService.ts https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/agent-core-v2/src/app/sessionManager/sessionManagerService.ts
12.7 kB
| 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<string, ISessionScopeHandle>(); | |
| private readonly owners = new Map<string, SessionLifecycleService>(); | |
| private readonly pendingResumes = new Map<string, Promise<ISessionScopeHandle | undefined>>(); | |
| private readonly resumeFailures = new Map<string, Error>(); | |
| private readonly lifecycleChains = new Map<string, Promise<void>>(); | |
| private readonly controllers = new Map<string, SessionControllerEntry>(); | |
| private readonly controllerEntries = new Set<SessionControllerEntry>(); | |
| private readonly willCreateEmitter = new Emitter<SessionWillCreateEvent>(); | |
| readonly onWillCreateSession: Event<SessionWillCreateEvent> = this.willCreateEmitter.event; | |
| private readonly didCreateEmitter = new Emitter<SessionCreatedEvent & IWaitUntil>(); | |
| readonly onDidCreateSession = this.didCreateEmitter.event; | |
| private readonly willCloseEmitter = new Emitter<SessionWillCloseEvent & IWaitUntil>(); | |
| readonly onWillCloseSession = this.willCloseEmitter.event; | |
| private readonly didCloseEmitter = new Emitter<SessionClosedEvent>(); | |
| readonly onDidCloseSession = this.didCloseEmitter.event; | |
| private readonly willDeleteEmitter = new Emitter<{ readonly sessionId: string } & IWaitUntil>(); | |
| readonly onWillDeleteSession = this.willDeleteEmitter.event; | |
| private readonly didArchiveEmitter = new Emitter<SessionArchivedEvent>(); | |
| readonly onDidArchiveSession = this.didArchiveEmitter.event; | |
| private readonly didForkEmitter = new Emitter<SessionForkedEvent>(); | |
| readonly onDidForkSession = this.didForkEmitter.event; | |
| constructor( | |
| private readonly workspaces: IWorkspaceInstanceManager, | |
| private readonly index: ISessionIndex, | |
| ) {} | |
| async create(options: CreateManagedSessionOptions): Promise<ISessionScopeHandle> { | |
| 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<ISessionScopeHandle | undefined> { | |
| 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<SessionSummary | undefined> { | |
| return this.index.get(sessionId); | |
| } | |
| async whenResumeSettled(sessionId: string): Promise<void> { | |
| 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<T>(sessionId: string, work: () => Promise<T>): Promise<T> { | |
| 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<T>(keys: readonly string[], work: () => Promise<T>): Promise<T> { | |
| 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<T>( | |
| sessionId: string, | |
| work: (unguarded: UnguardedSessionLifecycle) => Promise<T>, | |
| ): Promise<T> { | |
| 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<void> { | |
| await this.serializeLifecycle(sessionId, async () => this.owners.get(sessionId)?.close(sessionId)); | |
| } | |
| private async archiveInner(sessionId: string): Promise<void> { | |
| await (await this.controllerForSession(sessionId))?.archive(sessionId); | |
| } | |
| async archive(sessionId: string): Promise<void> { | |
| await this.serializeLifecycle(sessionId, () => this.archiveInner(sessionId)); | |
| } | |
| private async restoreInner( | |
| sessionId: string, | |
| options?: ResumeSessionOptions, | |
| ): Promise<ISessionScopeHandle | undefined> { | |
| return (await this.controllerForSession(sessionId))?.restore(sessionId, options); | |
| } | |
| async restore(sessionId: string, options?: ResumeSessionOptions): Promise<ISessionScopeHandle | undefined> { | |
| return this.serializeLifecycle(sessionId, () => this.restoreInner(sessionId, options)); | |
| } | |
| async delete(sessionId: string): Promise<void> { | |
| 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<unknown>[] = []; | |
| 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<SessionMeta> { | |
| 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<SessionMeta> { | |
| 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<SessionLifecycleService | undefined> { | |
| 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'); | |