diff --git a/.changeset/clean-session-init-failures.md b/.changeset/clean-session-init-failures.md new file mode 100644 index 0000000000..2a8903a095 --- /dev/null +++ b/.changeset/clean-session-init-failures.md @@ -0,0 +1,6 @@ +--- +"@moonshot-ai/agent-core-v2": patch +"@moonshot-ai/kimi-code": patch +--- + +Release failed session initializations, hide half-initialized handles, and block overlapping same-ID attempts from writing the same session files. diff --git a/packages/agent-core-v2/src/app/sessionIndex/sessionIndexService.ts b/packages/agent-core-v2/src/app/sessionIndex/sessionIndexService.ts index 00ac2f7e06..e00f17ab1d 100644 --- a/packages/agent-core-v2/src/app/sessionIndex/sessionIndexService.ts +++ b/packages/agent-core-v2/src/app/sessionIndex/sessionIndexService.ts @@ -21,8 +21,10 @@ * summary is resolved through the read model — falling back to a disk read + * backfill on a cold miss. Writes (create / archive / metadata update) keep the * read model warm via `SessionMetadata`; new sessions that have not been - * mirrored yet are simply a cold miss and backfilled on first read. The legacy - * N+1 path remains as the flag-off fallback — and as the runtime fallback when + * mirrored yet are simply a cold miss and backfilled on first read. Persisted + * metadata remains the membership source, so a cached row is ignored after its + * state document disappears. The legacy N+1 path remains as the flag-off + * fallback — and as the runtime fallback when * the query store reports `storage.locked` (another process holds the writer * lock): the first lock warns once and disables the read model for the rest of * the process lifetime. @@ -193,7 +195,14 @@ export class FileSessionIndex implements ISessionIndex { private async getFromReadModel(id: string): Promise { const cached = await this.queryStore.get(SESSION_COLLECTION, id); - if (cached !== undefined) return cached; + if (cached !== undefined && typeof cached.workspaceId === 'string') { + const summary = await this.readSummary(cached.workspaceId, id); + if (summary !== undefined) { + if (cached.id === id && cached.createdAt === summary.createdAt) return cached; + await this.queryStore.put(SESSION_COLLECTION, id, summary); + return summary; + } + } // Cold miss: locate the session on disk, then read + backfill. for (const workspaceId of await this.listWorkspaceIds()) { if (!(await this.hasSession(workspaceId, id))) continue; @@ -240,11 +249,17 @@ export class FileSessionIndex implements ISessionIndex { sessionId: string, ): Promise { const cached = await this.queryStore.get(SESSION_COLLECTION, sessionId); - if (cached !== undefined) return cached; const summary = await this.readSummary(workspaceId, sessionId); - if (summary !== undefined) { - await this.queryStore.put(SESSION_COLLECTION, sessionId, summary); + if (summary === undefined) return undefined; + if ( + cached !== undefined && + cached.id === sessionId && + cached.workspaceId === workspaceId && + cached.createdAt === summary.createdAt + ) { + return cached; } + await this.queryStore.put(SESSION_COLLECTION, sessionId, summary); return summary; } diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts index ff6178c39c..4369da57a8 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts @@ -107,26 +107,29 @@ export interface ISessionLifecycleService { create(opts: CreateSessionOptions): Promise; /** * Return the live handle for `sessionId`, or `undefined` when it is not open. - * A session whose cold {@link resume} is still in flight is intentionally NOT - * returned — its main agent has not finished restore + replay, so the handle - * is half-initialized. Callers that must obtain the handle should - * `await resume(sessionId)` instead. This invisibility is a service - * invariant, not caller discipline: every read path (`get` / {@link list} / - * {@link resume}) agrees a resuming session is not yet observable. + * A session whose create, fork, or cold {@link resume} initialization has not + * published its metadata, MCP readiness, and any required main-agent restore + * is intentionally NOT returned. Once that core state is ready, the handle + * is visible while creation hooks finish, unless an explicit close/archive + * has taken ownership of teardown. Callers that must wait for publication + * should `await resume(sessionId)` instead. */ get(sessionId: string): ISessionScopeHandle | undefined; /** - * Snapshot of every fully-initialized live session. Excludes sessions still - * mid-{@link resume} for the same reason as {@link get}. + * Snapshot of every published live session. Excludes sessions whose core + * initialization has not reached the publication point described by + * {@link get}. */ list(): readonly ISessionScopeHandle[]; /** * Load a persisted session into the live scope tree and restore its main - * agent from the persisted wire log. Returns the existing handle when the - * session is already live (a no-op in that case — live agents are never - * re-restored). Returns `undefined` when the session is unknown to the index - * or neither the persisted session summary nor the workspace registry can - * provide a workdir (mirrors the cold-source limitation of `fork`). + * agent from the persisted wire log. Returns the existing published handle + * when the session is already live, and waits for an unpublished same-id + * initialization to publish or fail before retrying the lookup. An + * initializing session already claimed by close/archive is unavailable. + * Returns `undefined` when the session is unknown to the index or neither the + * persisted session summary nor the workspace registry can provide a workdir + * (mirrors the cold-source limitation of `fork`). * * Lets the read edges (snapshot / messages) serve cold sessions — created by * a previous process or by v1 — without requiring a prior `create` in this diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts index 25a0e47492..2f6a4b4bd9 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -15,11 +15,18 @@ * (TUI, export) can discover sessions created by the v2 engine. Fork flushes * live agent logs and rejects non-empty logs without a protocol metadata * envelope instead of stamping legacy data as current. + * Failed attempts release only the Session handle they created. A failed fresh + * create also removes the session directory it exclusively claimed before + * materialization after queued metadata writes settle; resume failures preserve + * existing persisted state. Per-session initialization claims keep writers + * exclusive while publishing only core-ready handles and preserving explicit + * close/archive requests made before a handle exists. The same claim keeps an + * initializing handle private while explicit teardown is in progress. */ import { randomUUID } from 'node:crypto'; -import { join } from 'pathe'; +import { dirname, join } from 'pathe'; import { ulid } from 'ulid'; import { InstantiationType } from '#/_base/di/extensions'; @@ -95,11 +102,36 @@ import { type MaterializeSessionOptions = Omit & { readonly sessionId: string; readonly workspaceId?: string; + readonly claim: SessionInitializationClaim; + readonly forkDestination?: boolean; + readonly directoryOwnership?: SessionDirectoryOwnership; }; +interface SessionDirectoryOwnership { + sessionDir?: string; + owned: boolean; +} + +interface SessionInitializationClaim { + readonly published: boolean; + readonly closing: boolean; + readonly publishedOrSettled: Promise; + readonly settled: Promise; + markPublished(): void; + startLifecycleAction(): boolean; + finishLifecycleAction(): void; + release(): void; +}; + +interface SessionLifecycleAction { + readonly handle: ISessionScopeHandle; + readonly claim?: SessionInitializationClaim; +} + export class SessionLifecycleService extends Disposable implements ISessionLifecycleService { declare readonly _serviceBrand: undefined; private readonly sessions = new Map(); + private readonly initializations = new Map(); private readonly _onDidCreateSession = this._register(new Emitter()); readonly onDidCreateSession: Event = this._onDidCreateSession.event; private readonly _onDidCloseSession = this._register(new Emitter()); @@ -112,12 +144,6 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec 'onDidCreateSession', 'onWillCloseSession', ]); - /** In-flight `resume` promises, keyed by session id. De-dupes concurrent cold - * loads so a hot read path (e.g. snapshot retry) cannot materialize the same - * session twice and leak a handle — and doubles as the visibility gate for - * `get` / `list`: while an id is present here its materialized handle is - * half-initialized (main agent not yet restored + replayed) and must not be - * observable. */ private readonly resuming = new Map>(); constructor( @@ -141,14 +167,43 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec async create(opts: CreateSessionOptions): Promise { const sessionId = opts.sessionId ?? createSessionId(); - const handle = await this.materializeSession({ ...opts, sessionId }); - await this.appendSessionIndexEntry(sessionId, opts.workDir); - if (this.config.get(DEFAULT_PLAN_MODE_SECTION) === true) { - const main = await ensureMainAgent(handle); - await main.accessor.get(IAgentPlanService).enter(); + const claim = this.claimInitialization(sessionId); + const directoryOwnership: SessionDirectoryOwnership = { owned: false }; + let handle: ISessionScopeHandle | undefined; + try { + handle = await this.materializeSession({ + ...opts, + sessionId, + claim, + directoryOwnership, + }); + await this.appendSessionIndexEntry(sessionId, opts.workDir); + if (this.config.get(DEFAULT_PLAN_MODE_SECTION) === true) { + const main = await ensureMainAgent(handle); + await main.accessor.get(IAgentPlanService).enter(); + } + this.publishSessionHandle(sessionId, handle, claim); + await this.announceCreated({ sessionId, handle, source: 'startup' }, claim); + this.assertSessionHandleOwned(sessionId, handle, claim); + return handle; + } catch (error) { + const failedHandle = handle; + if (failedHandle !== undefined && !claim.published && !claim.closing) { + await throwAfterCleanup(error, () => + directoryOwnership.owned + ? this.rollbackOwnedSessionDirectory( + sessionId, + failedHandle, + directoryOwnership.sessionDir, + claim, + ) + : this.disposeFailedSession(sessionId, failedHandle), + ); + } + throw error; + } finally { + claim.release(); } - await this.announceCreated({ sessionId, handle, source: 'startup' }); - return handle; } private async materializeSession(opts: MaterializeSessionOptions): Promise { @@ -186,32 +241,59 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec // synchronously in their constructors, so the probe must have landed by // the time the first Session-scoped service is resolved. await this.hostEnv.ready; - const handle = createScopedChildHandle( - this.instantiation, - LifecycleScope.Session, - opts.sessionId, - { - extra: [...sessionContextSeed(ctx)], - }, - ) as ISessionScopeHandle; - // Construct the Session activity kernel eagerly so its lane is `restoring` - // for the whole materialize / replay window — edge commands that arrive - // before `markActive()` are rejected with `activity.session_rejected`. - handle.accessor.get(ISessionActivityKernel); - if (additionalDirs.length > 0) { - // De-duplication happens inside setAdditionalDirs (resolve + Set), - // matching v1's normalizeAdditionalDirs. - handle.accessor.get(ISessionWorkspaceContext).setAdditionalDirs(additionalDirs); + let handle: ISessionScopeHandle | undefined; + let directoryOwned = false; + try { + if (opts.directoryOwnership !== undefined || opts.forkDestination === true) { + await this.claimSessionDirectory(opts.sessionId, ctx.sessionDir); + directoryOwned = true; + if (opts.directoryOwnership !== undefined) { + opts.directoryOwnership.sessionDir = ctx.sessionDir; + opts.directoryOwnership.owned = true; + } + } + handle = createScopedChildHandle( + this.instantiation, + LifecycleScope.Session, + opts.sessionId, + { + extra: [...sessionContextSeed(ctx)], + }, + ) as ISessionScopeHandle; + handle.accessor.get(ISessionActivityKernel); + if (additionalDirs.length > 0) { + handle.accessor.get(ISessionWorkspaceContext).setAdditionalDirs(additionalDirs); + } + this.registerSessionHandle(opts.sessionId, handle, opts.claim); + await handle.accessor.get(ISessionMetadata).ready; + void handle.accessor.get(ISessionSkillCatalog).ready; + await handle.accessor.get(IAgentLifecycleService).ensureMcpReady(opts.mcpServers); + handle.accessor.get(ISessionExternalHooksService); + return handle; + } catch (error) { + const failedHandle = handle; + if (!opts.claim.closing) { + if (directoryOwned && failedHandle !== undefined) { + await throwAfterCleanup(error, () => + this.rollbackOwnedSessionDirectory( + opts.sessionId, + failedHandle, + ctx.sessionDir, + opts.claim, + ), + ); + } else if (directoryOwned) { + await throwAfterCleanup(error, () => + this.removeOwnedSessionDirectory(opts.sessionId, ctx.sessionDir, opts.claim), + ); + } else if (failedHandle !== undefined) { + await throwAfterCleanup(error, () => + this.disposeFailedSession(opts.sessionId, failedHandle), + ); + } + } + throw error; } - this.sessions.set(opts.sessionId, handle); - await handle.accessor.get(ISessionMetadata).ready; - void handle.accessor.get(ISessionSkillCatalog).ready; - // First `ensureMcpReady` call for the session — it starts the initial MCP - // load, so the caller-supplied servers must ride on it (later calls, e.g. - // from agent creation, only await the in-flight load). - await handle.accessor.get(IAgentLifecycleService).ensureMcpReady(opts.mcpServers); - handle.accessor.get(ISessionExternalHooksService); - return handle; } private async appendSessionIndexEntry(sessionId: string, workDir: string): Promise { @@ -225,119 +307,197 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await this.appendLogStore.flush(); } - private async announceCreated(event: SessionCreatedEvent): Promise { - await this.hooks.onDidCreateSession.run(event); + private async announceCreated( + event: SessionCreatedEvent, + claim: SessionInitializationClaim, + ): Promise { + let hookFailure: { readonly error: unknown } | undefined; + try { + await this.hooks.onDidCreateSession.run(event); + } catch (error) { + hookFailure = { error }; + } + this.assertSessionHandleOwned(event.sessionId, event.handle, claim); this._onDidCreateSession.fire(event); // Deliberately broader than v1: resumes also emit, with `resumed: true` — // the flag exists precisely to distinguish them (v1's resume path never // emitted despite the schema having the flag). this.telemetry.track2('session_started', { resumed: event.source === 'resume' }); - event.handle.accessor.get(ISessionActivityKernel).markActive(); + if (hookFailure !== undefined) throw hookFailure.error; } get(sessionId: string): ISessionScopeHandle | undefined { - // A session mid-resume is already materialized in `this.sessions` (so - // close/archive can still find it) but its main agent has not finished - // restore + replay — exposing it would hand callers a half-initialized - // handle. Hide it until `resume` settles; callers that need the handle - // should `await resume(sessionId)`. - if (this.resuming.has(sessionId)) return undefined; + const initialization = this.initializations.get(sessionId); + if ( + initialization !== undefined && + (!initialization.published || initialization.closing) + ) { + return undefined; + } return this.sessions.get(sessionId); } resume(sessionId: string): Promise { - // Check in-flight resumes FIRST: `materializeSession` adds the session to - // `this.sessions` before `doResume` finishes restore/replay, so a concurrent - // caller that checks `sessions` first would get a half-initialized handle - // whose main agent has no context. Checking `resuming` first ensures - // concurrent callers wait for the full resume (including restore + replay) - // to complete. - const inflight = this.resuming.get(sessionId); - if (inflight !== undefined) return inflight; + const initialization = this.initializations.get(sessionId); + if (initialization !== undefined) { + if (initialization.closing) return Promise.resolve(undefined); + if (initialization.published) { + const live = this.sessions.get(sessionId); + if (live !== undefined) return Promise.resolve(live); + return Promise.resolve(undefined); + } + return initialization.publishedOrSettled.then(() => this.resume(sessionId)); + } const live = this.sessions.get(sessionId); if (live !== undefined) return Promise.resolve(live); + const inflight = this.resuming.get(sessionId); + if (inflight !== undefined) return inflight; const promise = this.doResume(sessionId) .catch((error: unknown) => { this.telemetry.track2('session_load_failed', { reason: isError2(error) ? error.code : error instanceof Error ? error.name : 'unknown', }); throw error; - }) - .finally(() => this.resuming.delete(sessionId)); - this.resuming.set(sessionId, promise); - return promise; + }); + const tracked = promise.finally(() => { + if (this.resuming.get(sessionId) === tracked) { + this.resuming.delete(sessionId); + } + }); + this.resuming.set(sessionId, tracked); + return tracked; } private async doResume(sessionId: string): Promise { - // Re-check after the serialized entry: a prior `resume` for the same id may - // have already materialized the session while this call was queued. - const live = this.sessions.get(sessionId); - if (live !== undefined) return live; + while (true) { + const initialization = this.initializations.get(sessionId); + if (initialization !== undefined) { + if (initialization.closing) return undefined; + if (initialization.published) return this.sessions.get(sessionId); + await initialization.publishedOrSettled; + continue; + } - const summary = await this.index.get(sessionId); - if (summary === undefined) return undefined; - const workspace = - summary.cwd === undefined ? await this.workspaceRegistry.get(summary.workspaceId) : undefined; - const workDir = summary.cwd ?? workspace?.root; - if (workDir === undefined) return undefined; + const live = this.sessions.get(sessionId); + if (live !== undefined) return live; - const handle = await this.materializeSession({ - sessionId, - workDir, - workspaceId: summary.workspaceId, - }); - const agents = handle.accessor.get(IAgentLifecycleService); - if (agents.getHandle(MAIN_AGENT_ID) === undefined) { - const main = await ensureMainAgent(handle); - // Resolve context memory BEFORE restoring so its reducers are registered; - // otherwise the wire replay applies context records into a void and the - // restored transcript never lands in context memory. - main.accessor.get(IAgentContextMemoryService); - const mainWireRecord = main.accessor.get(IAgentWireRecordService); - await mainWireRecord.restore(); - const records = mainWireRecord.getRecords() as readonly PersistedRecord[]; - await main.accessor.get(IAgentWireService).replay(...records); + const summary = await this.index.get(sessionId); + const workspace = + summary !== undefined && summary.cwd === undefined + ? await this.workspaceRegistry.get(summary.workspaceId) + : undefined; + const workDir = summary?.cwd ?? workspace?.root; + + const racedInitialization = this.initializations.get(sessionId); + if (racedInitialization !== undefined) { + if (racedInitialization.closing) return undefined; + if (racedInitialization.published) return this.sessions.get(sessionId); + await racedInitialization.publishedOrSettled; + continue; + } + const racedLive = this.sessions.get(sessionId); + if (racedLive !== undefined) return racedLive; + if (summary === undefined || workDir === undefined) return undefined; + + let claim: SessionInitializationClaim; + try { + claim = this.claimInitialization(sessionId); + } catch (error) { + if (isError2(error) && error.code === ErrorCodes.SESSION_ALREADY_EXISTS) continue; + throw error; + } + let handle: ISessionScopeHandle | undefined; + try { + handle = await this.materializeSession({ + sessionId, + workDir, + workspaceId: summary.workspaceId, + claim, + }); + const agents = handle.accessor.get(IAgentLifecycleService); + if (agents.getHandle(MAIN_AGENT_ID) === undefined) { + const main = await ensureMainAgent(handle); + main.accessor.get(IAgentContextMemoryService); + const mainWireRecord = main.accessor.get(IAgentWireRecordService); + await mainWireRecord.restore(); + const records = mainWireRecord.getRecords() as readonly PersistedRecord[]; + await main.accessor.get(IAgentWireService).replay(...records); + } + this.publishSessionHandle(sessionId, handle, claim); + await this.announceCreated({ sessionId, handle, source: 'resume' }, claim); + this.assertSessionHandleOwned(sessionId, handle, claim); + return handle; + } catch (error) { + if (handle !== undefined && !claim.published && !claim.closing) { + await this.disposeFailedSession(sessionId, handle); + } + throw error; + } finally { + claim.release(); + } } - await this.announceCreated({ sessionId, handle, source: 'resume' }); - return handle; } list(): readonly ISessionScopeHandle[] { - // Exclude sessions still mid-resume for the same reason as `get`: the handle - // exists but is not yet restored, so it must not be observable. const ready: ISessionScopeHandle[] = []; for (const [id, handle] of this.sessions) { - if (!this.resuming.has(id)) ready.push(handle); + const initialization = this.initializations.get(id); + if ( + initialization === undefined || + (initialization.published && !initialization.closing) + ) { + ready.push(handle); + } } return ready; } async close(sessionId: string): Promise { - const handle = this.sessions.get(sessionId); - if (handle === undefined) return; - await this.announceWillClose({ sessionId, handle, reason: 'exit' }); - this.sessions.delete(sessionId); - handle.accessor.get(ISessionActivityKernel).beginClosing(); - await this.drainAgents(handle); - handle.dispose(); - this._onDidCloseSession.fire({ sessionId }); + const action = await this.handleForLifecycleAction(sessionId); + if (action === undefined) return; + const { handle } = action; + try { + await this.announceWillClose({ sessionId, handle, reason: 'exit' }); + this.sessions.delete(sessionId); + handle.accessor.get(ISessionActivityKernel).beginClosing(); + await this.drainAgents(handle); + handle.dispose(); + this._onDidCloseSession.fire({ sessionId }); + } catch (error) { + if (action.claim !== undefined) { + await this.disposeFailedSession(sessionId, handle); + } + throw error; + } finally { + action.claim?.finishLifecycleAction(); + } } async archive(sessionId: string): Promise { - const handle = this.sessions.get(sessionId); - if (handle === undefined) return; - const meta = handle.accessor.get(ISessionMetadata); - await meta.setArchived(true); - handle.accessor.get(ISessionActivityKernel).beginClosing(); - await this.drainAgents(handle); - this.event.publish({ - type: 'event.session.archived', - payload: { sessionId }, - }); - await this.announceWillClose({ sessionId, handle, reason: 'exit' }); - this.sessions.delete(sessionId); - handle.dispose(); - this._onDidArchiveSession.fire({ sessionId }); + const action = await this.handleForLifecycleAction(sessionId); + if (action === undefined) return; + const { handle } = action; + try { + const meta = handle.accessor.get(ISessionMetadata); + await meta.setArchived(true); + handle.accessor.get(ISessionActivityKernel).beginClosing(); + await this.drainAgents(handle); + this.event.publish({ + type: 'event.session.archived', + payload: { sessionId }, + }); + await this.announceWillClose({ sessionId, handle, reason: 'exit' }); + this.sessions.delete(sessionId); + handle.dispose(); + this._onDidArchiveSession.fire({ sessionId }); + } catch (error) { + if (action.claim !== undefined) { + await this.disposeFailedSession(sessionId, handle); + } + throw error; + } finally { + action.claim?.finishLifecycleAction(); + } } async restore(sessionId: string): Promise { @@ -358,12 +518,219 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } } + private async disposeFailedSession( + sessionId: string, + handle: ISessionScopeHandle, + ): Promise { + if (this.sessions.get(sessionId) === handle) { + this.sessions.delete(sessionId); + } + try { + handle.accessor.get(ISessionActivityKernel).beginClosing(); + } catch {} + try { + const agentLifecycle = handle.accessor.get(IAgentLifecycleService); + for (const agent of agentLifecycle.list()) { + try { + await agentLifecycle.remove(agent.id); + } catch {} + } + } catch {} + try { + await handle.accessor.get(ISessionMetadata).whenIdle(); + } catch {} + try { + handle.dispose(); + } catch {} + } + + private claimInitialization(sessionId: string): SessionInitializationClaim { + if (this.sessions.has(sessionId) || this.initializations.has(sessionId)) { + throw sessionAlreadyExistsError(sessionId); + } + let writerReleased = false; + let published = false; + let lifecycleActionStarted = false; + let lifecycleActionFinished = false; + let claimSettled = false; + let resolvePublishedOrSettled!: () => void; + let settle!: () => void; + const publishedOrSettled = new Promise((resolve) => { + resolvePublishedOrSettled = resolve; + }); + const settled = new Promise((resolve) => { + settle = resolve; + }); + const trySettle = (): void => { + if ( + claimSettled || + !writerReleased || + (lifecycleActionStarted && !lifecycleActionFinished) + ) { + return; + } + claimSettled = true; + if (this.initializations.get(sessionId) === claim) { + this.initializations.delete(sessionId); + } + resolvePublishedOrSettled(); + settle(); + }; + const claim: SessionInitializationClaim = { + get published() { + return published; + }, + get closing() { + return lifecycleActionStarted && published; + }, + publishedOrSettled, + settled, + markPublished: () => { + if (writerReleased || published) return; + published = true; + resolvePublishedOrSettled(); + }, + startLifecycleAction: () => { + if (claimSettled || writerReleased || lifecycleActionStarted) { + return false; + } + lifecycleActionStarted = true; + return true; + }, + finishLifecycleAction: () => { + if (!lifecycleActionStarted || lifecycleActionFinished) return; + lifecycleActionFinished = true; + trySettle(); + }, + release: () => { + if (writerReleased) return; + writerReleased = true; + resolvePublishedOrSettled(); + trySettle(); + }, + }; + this.initializations.set(sessionId, claim); + return claim; + } + + private registerSessionHandle( + sessionId: string, + handle: ISessionScopeHandle, + claim: SessionInitializationClaim, + ): void { + if (this.initializations.get(sessionId) !== claim || this.sessions.has(sessionId)) { + throw sessionAlreadyExistsError(sessionId); + } + this.sessions.set(sessionId, handle); + } + + private publishSessionHandle( + sessionId: string, + handle: ISessionScopeHandle, + claim: SessionInitializationClaim, + ): void { + this.assertSessionHandleOwned(sessionId, handle, claim); + handle.accessor.get(ISessionActivityKernel).markActive(); + claim.markPublished(); + this.assertSessionHandleOwned(sessionId, handle, claim); + } + + private assertSessionHandleOwned( + sessionId: string, + handle: ISessionScopeHandle, + claim: SessionInitializationClaim, + ): void { + if ( + this.initializations.get(sessionId) !== claim || + this.sessions.get(sessionId) !== handle || + claim.closing + ) { + throw new Error2( + ErrorCodes.SESSION_CLOSED, + `Session "${sessionId}" closed during initialization`, + ); + } + } + + private async handleForLifecycleAction( + sessionId: string, + ): Promise { + const initialization = this.initializations.get(sessionId); + if (initialization === undefined) { + const handle = this.sessions.get(sessionId); + return handle === undefined ? undefined : { handle }; + } + if (!initialization.startLifecycleAction()) return undefined; + if (!initialization.published) await initialization.publishedOrSettled; + const handle = this.sessions.get(sessionId); + if ( + this.initializations.get(sessionId) !== initialization || + !initialization.published || + handle === undefined + ) { + initialization.finishLifecycleAction(); + return undefined; + } + return { handle, claim: initialization }; + } + + private async forkSourceHandle(sessionId: string): Promise { + let initialization = this.initializations.get(sessionId); + while (initialization !== undefined) { + if (initialization.published && !initialization.closing) { + return this.sessions.get(sessionId); + } + await (initialization.closing + ? initialization.settled + : initialization.publishedOrSettled); + initialization = this.initializations.get(sessionId); + } + return this.sessions.get(sessionId); + } + + private async claimSessionDirectory(sessionId: string, sessionDir: string): Promise { + await this.hostFs.mkdir(dirname(sessionDir), { recursive: true }); + try { + await this.hostFs.mkdir(sessionDir); + } catch (error) { + if (isExistingFileError(error)) throw sessionAlreadyExistsError(sessionId); + throw error; + } + } + + private async rollbackOwnedSessionDirectory( + sessionId: string, + handle: ISessionScopeHandle, + sessionDir: string | undefined, + claim: SessionInitializationClaim, + ): Promise { + if (this.initializations.get(sessionId) !== claim || sessionDir === undefined) { + await this.disposeFailedSession(sessionId, handle); + return; + } + if (this.sessions.get(sessionId) === handle) { + this.sessions.delete(sessionId); + } + await this.disposeFailedSession(sessionId, handle); + await this.appendLogStore.flush(); + await this.removeOwnedSessionDirectory(sessionId, sessionDir, claim); + } + + private async removeOwnedSessionDirectory( + sessionId: string, + sessionDir: string | undefined, + claim: SessionInitializationClaim, + ): Promise { + if (this.initializations.get(sessionId) !== claim || sessionDir === undefined) return; + await this.hostFs.remove(sessionDir); + } + async fork(opts: ForkSessionOptions): Promise { const sourceId = opts.sourceSessionId; // 1. Resolve the source: prefer a live handle, otherwise fall back to the // persisted index (so a closed session can still be forked, like v1). - const sourceHandle = this.sessions.get(sourceId); + const sourceHandle = await this.forkSourceHandle(sourceId); const indexSummary = await this.index.get(sourceId); if (sourceHandle === undefined && indexSummary === undefined) { throw new Error2(ErrorCodes.SESSION_NOT_FOUND, `session ${sourceId} does not exist`); @@ -384,6 +751,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec let targetId: string | undefined; let target: ISessionScopeHandle | undefined; let targetSessionDir: string | undefined; + let targetClaim: SessionInitializationClaim | undefined; try { // 3. Resolve the work dir the fork inherits (same workspace as the source). const workspace = await this.workspaceRegistry.get(workspaceId); @@ -400,16 +768,16 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec // 5. Mint the target id and reject collisions. targetId = opts.newSessionId ?? createSessionId(); if (this.sessions.has(targetId) || (await this.index.get(targetId)) !== undefined) { - throw new Error2( - ErrorCodes.SESSION_ALREADY_EXISTS, - `Session "${targetId}" already exists`, - ); + throw sessionAlreadyExistsError(targetId); } + targetClaim = this.claimInitialization(targetId); // 6. Materialize the target session scope (fresh metadata + storage). target = await this.materializeSession({ sessionId: targetId, workDir: workspace.root, + claim: targetClaim, + forkDestination: true, }); const targetCtx = target.accessor.get(ISessionContext); targetSessionDir = targetCtx.sessionDir; @@ -473,33 +841,41 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } await this.appendSessionIndexEntry(targetId, workspace.root); + this.publishSessionHandle(targetId, target, targetClaim); this._onDidForkSession.fire({ sourceSessionId: sourceId, sessionId: targetId, handle: target, }); - await this.announceCreated({ sessionId: targetId, handle: target, source: 'fork' }); + await this.announceCreated( + { sessionId: targetId, handle: target, source: 'fork' }, + targetClaim, + ); + this.assertSessionHandleOwned(targetId, target, targetClaim); return target; } catch (error) { - // Roll back the half-fork, mirroring v1's `rm -rf` of the target dir: - // drop the materialized handle from the live registry (otherwise a - // retry with the same id trips SESSION_ALREADY_EXISTS on the in-memory - // check) and delete whatever was copied to disk. - if (targetId !== undefined) { - this.sessions.delete(targetId); - } - if (target !== undefined) { - try { - target.dispose(); - } catch { - // best effort — the session dir is removed below regardless - } - } - if (targetSessionDir !== undefined) { - await this.hostFs.remove(targetSessionDir).catch(() => {}); + if ( + targetId !== undefined && + target !== undefined && + targetClaim !== undefined && + !targetClaim.published && + !targetClaim.closing + ) { + const failedTargetId = targetId; + const failedTarget = target; + const failedTargetClaim = targetClaim; + await throwAfterCleanup(error, () => + this.rollbackOwnedSessionDirectory( + failedTargetId, + failedTarget, + targetSessionDir, + failedTargetClaim, + ), + ); } throw error; } finally { + targetClaim?.release(); quiesce?.dispose(); } } @@ -702,6 +1078,28 @@ function isMissingFileError(error: unknown): boolean { return code === 'ENOENT'; } +function isExistingFileError(error: unknown): boolean { + const unwrapped = unwrapErrorCause(error); + if (unwrapped === null || typeof unwrapped !== 'object') return false; + const code = (unwrapped as { readonly code?: unknown }).code; + return code === 'EEXIST'; +} + +async function throwAfterCleanup( + error: unknown, + cleanup: () => Promise, +): Promise { + try { + await cleanup(); + } catch (cleanupError) { + throw new AggregateError( + [error, cleanupError], + 'Session initialization and cleanup both failed', + ); + } + throw error; +} + /** * Mint a session id in the canonical `session_` form, matching * v1's `createSessionId` (`packages/agent-core/src/rpc/core-impl.ts`). @@ -714,6 +1112,13 @@ function createSessionId(): string { return `session_${randomUUID()}`; } +function sessionAlreadyExistsError(sessionId: string): Error2 { + return new Error2( + ErrorCodes.SESSION_ALREADY_EXISTS, + `Session "${sessionId}" already exists`, + ); +} + function freshMetadataRecord(): PersistedWireRecord { return { type: 'metadata', diff --git a/packages/agent-core-v2/src/session/sessionMetadata/sessionMetadata.ts b/packages/agent-core-v2/src/session/sessionMetadata/sessionMetadata.ts index be5fadf508..7c8e0c396b 100644 --- a/packages/agent-core-v2/src/session/sessionMetadata/sessionMetadata.ts +++ b/packages/agent-core-v2/src/session/sessionMetadata/sessionMetadata.ts @@ -5,8 +5,9 @@ * layers to read and update the session's durable metadata (title, timestamps, * archived flag, fork provenance). Owns the in-memory copy, persists it as a * single atomic document through `storage`, and notifies changes via - * `onDidChangeMetadata`. Session-scoped — one instance per session. The initial - * document is materialized when the session is created. + * `onDidChangeMetadata`. Exposes an idle barrier so session teardown can wait + * for queued document and read-model writes. Session-scoped — one instance per + * session. The initial document is materialized when the session is created. */ import type { Event } from '#/_base/event'; @@ -78,6 +79,7 @@ export interface ISessionMetadata { readonly ready: Promise; readonly onDidChangeMetadata: Event; + whenIdle(): Promise; read(): Promise; update(patch: SessionMetaPatch): Promise; setTitle(title: string): Promise; diff --git a/packages/agent-core-v2/src/session/sessionMetadata/sessionMetadataService.ts b/packages/agent-core-v2/src/session/sessionMetadata/sessionMetadataService.ts index c4aaef27d5..e72ea03b68 100644 --- a/packages/agent-core-v2/src/session/sessionMetadata/sessionMetadataService.ts +++ b/packages/agent-core-v2/src/session/sessionMetadata/sessionMetadataService.ts @@ -11,9 +11,11 @@ * update is persisted, the fresh summary is mirrored into the `IQueryStore` * derived read model so `FileSessionIndex` can serve listings without * re-reading `state.json`. Mirroring is best-effort (a failure is logged, not - * thrown) and is a no-op when the flag is off. Initial creation in `load()` is - * intentionally not mirrored — a not-yet-mirrored session is simply a cold - * read-model miss that `FileSessionIndex` backfills on first read. + * thrown) and is a no-op when the flag is off. The metadata owner exposes an + * idle barrier covering queued persistence and mirroring before Session-scope + * teardown. Initial creation in `load()` is intentionally not mirrored — a + * not-yet-mirrored session is simply a cold read-model miss that + * `FileSessionIndex` backfills on first read. */ import { InstantiationType } from '#/_base/di/extensions'; @@ -69,6 +71,15 @@ export class SessionMetadata extends Disposable implements ISessionMetadata { return this.data; } + async whenIdle(): Promise { + await this.ready; + let queued: Promise; + do { + queued = this.updateQueue; + await queued; + } while (this.updateQueue !== queued); + } + async update(patch: SessionMetaPatch): Promise { return this.enqueueUpdate(() => this.applyUpdate(patch)); } diff --git a/packages/agent-core-v2/test/app/sessionExport/sessionExport.test.ts b/packages/agent-core-v2/test/app/sessionExport/sessionExport.test.ts index c55a773e63..8841d11b7c 100644 --- a/packages/agent-core-v2/test/app/sessionExport/sessionExport.test.ts +++ b/packages/agent-core-v2/test/app/sessionExport/sessionExport.test.ts @@ -403,6 +403,7 @@ function stubSessionMetadata(meta: SessionMeta): ISessionMetadata { _serviceBrand: undefined, ready: Promise.resolve(), onDidChangeMetadata: noopEvent, + whenIdle: async () => {}, read: async () => meta, update: async () => {}, setTitle: async () => {}, diff --git a/packages/agent-core-v2/test/app/sessionIndex/sessionIndex.test.ts b/packages/agent-core-v2/test/app/sessionIndex/sessionIndex.test.ts index cd2bc1c719..5719d17d9b 100644 --- a/packages/agent-core-v2/test/app/sessionIndex/sessionIndex.test.ts +++ b/packages/agent-core-v2/test/app/sessionIndex/sessionIndex.test.ts @@ -283,13 +283,52 @@ describe('FileSessionIndex (read model)', () => { }); it('get prefers the read model over disk', async () => { + await seedSession('warm', { title: 'from disk', createdAt: 1, updatedAt: 2 }); const store = build(); - // Not seeded on disk — only present in the read model. await queryStore.put(SESSION_COLLECTION, 'warm', summary('warm', { title: 'cached' })); const got = await store.get('warm'); expect(got?.title).toBe('cached'); }); + it('does not expose a cached session without persisted state', async () => { + await fsp.mkdir(join(sessionsDir, workspaceId, 'initializing'), { recursive: true }); + const store = build(); + await queryStore.put( + SESSION_COLLECTION, + 'initializing', + summary('initializing', { title: 'cached' }), + ); + + await expect(store.get('initializing')).resolves.toBeUndefined(); + await expect( + store.list({ sessionId: 'initializing', includeArchived: true }), + ).resolves.toEqual({ items: [] }); + }); + + it('replaces a cached summary when the same path contains a new session', async () => { + await seedSession('recreated', { + title: 'current', + createdAt: 10, + updatedAt: 11, + }); + const store = build(); + await queryStore.put( + SESSION_COLLECTION, + 'recreated', + summary('recreated', { title: 'stale', createdAt: 1, updatedAt: 2 }), + ); + + await expect(store.get('recreated')).resolves.toMatchObject({ + id: 'recreated', + title: 'current', + createdAt: 10, + }); + await expect(queryStore.get(SESSION_COLLECTION, 'recreated')).resolves.toMatchObject({ + title: 'current', + createdAt: 10, + }); + }); + it('list filters by childOf from the read model', async () => { await seedSession('child-a', { createdAt: 2, @@ -319,8 +358,8 @@ describe('FileSessionIndex (read model)', () => { }); it('countActive reflects read-model updates', async () => { - await seedSession('a', {}); - await seedSession('b', { archived: true }); + await seedSession('a', { createdAt: 1, updatedAt: 2 }); + await seedSession('b', { archived: true, createdAt: 1, updatedAt: 2 }); const store = build(); expect(await store.countActive(workspaceId)).toBe(1); diff --git a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts index ec0241861d..2a33dc125e 100644 --- a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts +++ b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts @@ -1,3 +1,13 @@ +/** + * `sessionLifecycle` domain (L6) — App-to-Session scope lifecycle scenarios. + * + * Covers public create, resume, close, archive, fork, and failed-initialization + * recovery contracts through the real scoped host. Persistence, host, config, + * and agent boundaries are controlled through their service interfaces. + * + * Run: pnpm --filter @moonshot-ai/agent-core-v2 exec vitest run test/app/sessionLifecycle/sessionLifecycle.test.ts + */ + import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { mkdtemp, mkdir, readFile, rm, stat, writeFile } from 'node:fs/promises'; @@ -5,10 +15,14 @@ import { tmpdir } from 'node:os'; import { isAbsolute, join, resolve } from 'node:path'; import { InstantiationType } from '#/_base/di/extensions'; +import { SyncDescriptor } from '#/_base/di/descriptors'; import { Disposable } from '#/_base/di/lifecycle'; +import { ILogService } from '#/_base/log/log'; import { type IAgentScopeHandle, + type ISessionScopeHandle, LifecycleScope, + type ScopeSeed, _clearScopedRegistryForTests, registerScopedService, } from '#/_base/di/scope'; @@ -16,6 +30,7 @@ import { type ScopedTestHost, createScopedTestHost, stubPair } from '#/_base/di/ import { Event } from '#/_base/event'; import { IBootstrapService } from '#/app/bootstrap/bootstrap'; import { IConfigService } from '#/app/config/config'; +import { IFlagService } from '#/app/flag/flag'; import { IHostEnvironment } from '#/os/interface/hostEnvironment'; import { IHostFileSystem } from '#/os/interface/hostFileSystem'; import { HostFileSystem } from '#/os/backends/node-local/hostFsService'; @@ -40,8 +55,13 @@ import { ISessionActivity } from '#/session/sessionActivity/sessionActivity'; import { ISessionMetadata } from '#/session/sessionMetadata/sessionMetadata'; import { ISessionSkillCatalog } from '#/session/sessionSkillCatalog/skillCatalog'; import { ISessionIndex, type SessionSummary } from '#/app/sessionIndex/sessionIndex'; +import { FileSessionIndex } from '#/app/sessionIndex/sessionIndexService'; +import { JsonAtomicDocumentStore } from '#/persistence/backends/node-fs/atomicDocumentStore'; +import { FileStorageService } from '#/persistence/backends/node-fs/fileStorageService'; import { IAppendLogStore } from '#/persistence/interface/appendLogStore'; import { IAtomicDocumentStore } from '#/persistence/interface/atomicDocumentStore'; +import { IQueryStore } from '#/persistence/interface/queryStore'; +import { IFileSystemStorageService } from '#/persistence/interface/storage'; import { IWorkspaceLocalConfigService } from '#/app/workspaceLocalConfig/workspaceLocalConfig'; import { ISessionWorkspaceContext } from '#/session/workspaceContext/workspaceContext'; import { SessionWorkspaceContextService } from '#/session/workspaceContext/workspaceContextService'; @@ -50,24 +70,22 @@ import { encodeWorkDirKey } from '#/_base/utils/workdir-slug'; import { ISessionContext } from '#/session/sessionContext/sessionContext'; import { ITelemetryService } from '#/app/telemetry/telemetry'; import { Error2, ErrorCodes } from '#/errors'; +import { SessionMetadata } from '#/session/sessionMetadata/sessionMetadataService'; +import { stubLog } from '../../_base/log/stubs'; import { stubSessionActivityKernel } from '../../activity/stubs'; import { recordingTelemetry, type TelemetryRecord } from '../../app/telemetry/stubs'; +import { stubFlag } from '../flag/stubs'; +import { stubQueryStore } from '../../persistence/interface/stubs'; -function bootstrapStub(): IBootstrapService { - return { - sessionsDir: '/tmp/sessions', - homeDir: '/tmp', - sessionScope: (workspaceId: string, sessionId: string) => - `sessions/${workspaceId}/${sessionId}`, - sessionDir: (workspaceId: string, sessionId: string) => - `/tmp/sessions/${workspaceId}/${sessionId}`, - } as IBootstrapService; +function bootstrapStub(root: string): IBootstrapService { + return tmpBootstrapStub(root); } function tmpBootstrapStub(root: string): IBootstrapService { return { sessionsDir: join(root, 'sessions'), homeDir: root, + scope: (name: string) => name, sessionScope: (workspaceId: string, sessionId: string) => `sessions/${workspaceId}/${sessionId}`, sessionDir: (workspaceId: string, sessionId: string) => @@ -102,6 +120,7 @@ function metadataStub(): ISessionMetadata { _serviceBrand: undefined, ready: Promise.resolve(), onDidChangeMetadata: () => ({ dispose: () => {} }), + whenIdle: () => Promise.resolve(), read: () => Promise.resolve({} as never), update: () => Promise.resolve(), setTitle: () => Promise.resolve(), @@ -258,6 +277,25 @@ function appendLogStoreStub(): IAppendLogStore { }; } +function statefulQueryStore(): IQueryStore { + const base = stubQueryStore(); + const values = new Map(); + const entryKey = (collection: string, key: string): string => `${collection}\u0000${key}`; + return { + ...base, + put: (collection, key, value) => { + values.set(entryKey(collection, key), value); + return Promise.resolve(); + }, + delete: (collection, key) => { + values.delete(entryKey(collection, key)); + return Promise.resolve(); + }, + get: (collection: string, key: string) => + Promise.resolve(values.get(entryKey(collection, key)) as T | undefined), + }; +} + function atomicDocumentStoreStub(): IAtomicDocumentStore { return { _serviceBrand: undefined, @@ -358,8 +396,18 @@ function agentLifecycleCapturingPlanSpy(opts: { mainPreexists?: boolean } = {}): return { lifecycle, enter, create }; } -function tick(): Promise { - return new Promise((resolve) => setTimeout(resolve, 0)); +function deferred(): { + readonly promise: Promise; + readonly resolve: (value: T) => void; + readonly reject: (reason?: unknown) => void; +} { + let resolve!: (value: T) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + return { promise, resolve, reject }; } class NoopSessionExternalHooksService implements ISessionExternalHooksService { @@ -395,11 +443,14 @@ describe('SessionLifecycleService', () => { let host: ScopedTestHost | undefined; let telemetryRecords: TelemetryRecord[]; let tmpRoots: string[]; + let defaultRoot: string; - beforeEach(() => { + beforeEach(async () => { recordedSessionHookEvents = []; telemetryRecords = []; tmpRoots = []; + defaultRoot = await mkdtemp(join(tmpdir(), 'kimi-session-lifecycle-test-')); + tmpRoots.push(defaultRoot); _clearScopedRegistryForTests(); registerScopedService( LifecycleScope.App, @@ -440,9 +491,9 @@ describe('SessionLifecycleService', () => { await Promise.all(tmpRoots.map((root) => rm(root, { recursive: true, force: true }))); }); - function build(extra: ReturnType[] = []): ISessionLifecycleService { + function build(extra: ScopeSeed = []): ISessionLifecycleService { host = createScopedTestHost([ - stubPair(IBootstrapService, bootstrapStub()), + stubPair(IBootstrapService, bootstrapStub(defaultRoot)), stubPair(ISessionMetadata, metadataStub()), stubPair(IHostEnvironment, hostEnvironmentStub()), stubPair(ISessionSkillCatalog, skillCatalogStub()), @@ -450,9 +501,11 @@ describe('SessionLifecycleService', () => { stubPair(ISessionIndex, sessionIndexStub()), stubPair(IAppendLogStore, appendLogStoreStub()), stubPair(IAtomicDocumentStore, atomicDocumentStoreStub()), + stubPair(IQueryStore, stubQueryStore()), stubPair(IEventService, eventStub()), stubPair(IAgentLifecycleService, agentLifecycleStub()), stubPair(IConfigService, configStub()), + stubPair(IFlagService, stubFlag(false)), stubPair(ISessionCronService, { _serviceBrand: undefined } as unknown as ISessionCronService), stubPair(ISessionActivityKernel, stubSessionActivityKernel()), stubPair(IWorkspaceLocalConfigService, workspaceLocalConfigStub()), @@ -518,7 +571,7 @@ describe('SessionLifecycleService', () => { key: 'session_index.jsonl', record: { sessionId: 's1', - sessionDir: `/tmp/sessions/${workspaceId}/s1`, + sessionDir: join(defaultRoot, 'sessions', workspaceId, 's1'), workDir: '/tmp/proj', }, }, @@ -607,7 +660,7 @@ describe('SessionLifecycleService', () => { expect(ctx?.cwd).toBe(workDir); expect(ctx?.workspaceId).toBe(indexedWorkspaceId); - expect(ctx?.sessionDir).toBe(`/tmp/sessions/${indexedWorkspaceId}/s1`); + expect(ctx?.sessionDir).toBe(join(defaultRoot, 'sessions', indexedWorkspaceId, 's1')); }); it('archive flags metadata, removes agents, publishes the event, and disposes the session', async () => { @@ -755,32 +808,876 @@ describe('SessionLifecycleService', () => { expect(recordedSessionHookEvents).toEqual(['create:startup:s1', 'close:exit:s1']); }); - it('waits for MCP initialization before create returns', async () => { - let resolveMcpReady: (() => void) | undefined; - const mcpReady = new Promise((resolve) => { - resolveMcpReady = resolve; + it('publishes a created handle to read APIs before the create hook completes', async () => { + const svc = build(); + let resumed: Awaited>; + let fetched: ReturnType; + let listed: ReturnType = []; + const hook = svc.hooks.onDidCreateSession.register('read-live-session', async (event, next) => { + resumed = await svc.resume(event.sessionId); + fetched = svc.get(event.sessionId); + listed = svc.list(); + await next(); + }); + + const created = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + + expect(resumed!).toBe(created); + expect(fetched!).toBe(created); + expect(listed).toEqual([created]); + hook.dispose(); + }); + + it('keeps a handle returned by resume when the create hook later fails', async () => { + const hookStarted = deferred(); + const finishHook = deferred(); + const failure = new Error('session creation hook failed'); + const createdEvents: ISessionScopeHandle[] = []; + const svc = build(); + svc.onDidCreateSession((event) => createdEvents.push(event.handle)); + const hook = svc.hooks.onDidCreateSession.register('fail-created-session', async () => { + hookStarted.resolve(); + await finishHook.promise; + throw failure; + }); + + const creating = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + const createResult = creating.then( + () => undefined, + (error: unknown) => error, + ); + await hookStarted.promise; + const resumed = await svc.resume('s1'); + + expect(resumed).toBeDefined(); + expect(resumed!.accessor.get(ISessionActivityKernel).lane()).toBe('active'); + finishHook.resolve(); + await expect(createResult).resolves.toBe(failure); + + expect(svc.get('s1')).toBe(resumed); + await expect(svc.resume('s1')).resolves.toBe(resumed); + expect(createdEvents).toEqual([resumed]); + hook.dispose(); + }); + + it('rejects create when its creation hook closes the published session', async () => { + const svc = build(); + const closed: string[] = []; + svc.onDidCloseSession((event) => closed.push(event.sessionId)); + const hook = svc.hooks.onDidCreateSession.register( + 'close-created-session', + async (event, next) => { + await svc.close(event.sessionId); + await next(); + }, + ); + + await expect( + svc.create({ sessionId: 's1', workDir: '/tmp/proj' }), + ).rejects.toMatchObject({ code: ErrorCodes.SESSION_CLOSED }); + + expect(closed).toEqual(['s1']); + expect(svc.get('s1')).toBeUndefined(); + hook.dispose(); + }); + + it('does not share a disposed handle with a second resume during the creation hook', async () => { + const hookStarted = deferred(); + const finishHook = deferred(); + const svc = build([ + stubPair(ISessionIndex, sessionIndexWithSummary('s1', '/tmp/proj')), + stubPair(IAgentLifecycleService, agentLifecycleWithMainStub()), + ]); + const hook = svc.hooks.onDidCreateSession.register( + 'hold-created-session', + async (_event, next) => { + hookStarted.resolve(); + await finishHook.promise; + await next(); + }, + ); + + const firstResume = svc.resume('s1'); + const firstRejected = expect(firstResume).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, + }); + await hookStarted.promise; + await svc.close('s1'); + const secondResume = svc.resume('s1'); + finishHook.resolve(); + + await expect(secondResume).resolves.toBeUndefined(); + await firstRejected; + expect(svc.get('s1')).toBeUndefined(); + hook.dispose(); + }); + + it('hides a session while close waits for its lifecycle hook', async () => { + const createHookStarted = deferred(); + const finishCreateHook = deferred(); + const closeHookStarted = deferred(); + const finishCloseHook = deferred(); + const svc = build(); + const createHook = svc.hooks.onDidCreateSession.register( + 'hold-created-session', + async (_event, next) => { + createHookStarted.resolve(); + await finishCreateHook.promise; + await next(); + }, + ); + const closeHook = svc.hooks.onWillCloseSession.register( + 'hold-closing-session', + async (_event, next) => { + closeHookStarted.resolve(); + await finishCloseHook.promise; + await next(); + }, + ); + const closed: string[] = []; + svc.onDidCloseSession((event) => closed.push(event.sessionId)); + + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + const createRejected = expect(create).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, + }); + await createHookStarted.promise; + const close = svc.close('s1'); + await closeHookStarted.promise; + finishCreateHook.resolve(); + await createRejected; + + expect(svc.get('s1')).toBeUndefined(); + expect(svc.list()).toEqual([]); + await expect(svc.resume('s1')).resolves.toBeUndefined(); + expect(closed).toEqual([]); + + finishCloseHook.resolve(); + await close; + expect(closed).toEqual(['s1']); + createHook.dispose(); + closeHook.dispose(); + }); + + it('does not expose an initializing session after its close hook fails', async () => { + const createHookStarted = deferred(); + const finishCreateHook = deferred(); + const failure = new Error('close hook failed'); + const svc = build(); + const createHook = svc.hooks.onDidCreateSession.register( + 'hold-created-session', + async (_event, next) => { + createHookStarted.resolve(); + await finishCreateHook.promise; + await next(); + }, + ); + const closeHook = svc.hooks.onWillCloseSession.register('fail-close', () => { + throw failure; + }); + + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + const createRejected = expect(create).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, + }); + await createHookStarted.promise; + await expect(svc.close('s1')).rejects.toBe(failure); + + expect(svc.get('s1')).toBeUndefined(); + expect(svc.list()).toEqual([]); + await expect(svc.resume('s1')).resolves.toBeUndefined(); + + finishCreateHook.resolve(); + await createRejected; + expect(svc.get('s1')).toBeUndefined(); + createHook.dispose(); + closeHook.dispose(); + }); + + it('ignores close while archive owns initialization shutdown', async () => { + const createHookStarted = deferred(); + const finishCreateHook = deferred(); + const archiveStarted = deferred(); + const finishArchive = deferred(); + const archivedValues: boolean[] = []; + const svc = build([ + stubPair(ISessionMetadata, { + ...metadataStub(), + setArchived: (value) => { + archivedValues.push(value); + archiveStarted.resolve(); + return finishArchive.promise; + }, + }), + ]); + const createHook = svc.hooks.onDidCreateSession.register( + 'hold-created-session', + async (_event, next) => { + createHookStarted.resolve(); + await finishCreateHook.promise; + await next(); + }, + ); + const closed: string[] = []; + const archived: string[] = []; + svc.onDidCloseSession((event) => closed.push(event.sessionId)); + svc.onDidArchiveSession((event) => archived.push(event.sessionId)); + + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + const createRejected = expect(create).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, }); + await createHookStarted.promise; + const archive = svc.archive('s1'); + await archiveStarted.promise; + await svc.close('s1'); + finishCreateHook.resolve(); + await createRejected; + + expect(archivedValues).toEqual([true]); + expect(closed).toEqual([]); + + finishArchive.resolve(); + await archive; + expect(archived).toEqual(['s1']); + createHook.dispose(); + }); + + it('keeps an initializing create private until it is published', async () => { + const mcpStarted = deferred(); + const mcpReady = deferred(); const svc = build([ stubPair(IAgentLifecycleService, { ...agentLifecycleStub(), - ensureMcpReady: () => mcpReady, + ensureMcpReady: () => { + mcpStarted.resolve(); + return mcpReady.promise; + }, }), ]); + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await mcpStarted.promise; + const resume = svc.resume('s1'); + + expect(svc.get('s1')).toBeUndefined(); + expect(svc.list()).toEqual([]); + + mcpReady.resolve(); + const resumed = await resume; + const created = await create; + expect(resumed).toBe(created); + expect(svc.get('s1')).toBe(created); + expect(svc.list()).toEqual([created]); + }); + + it('retries resume lookup after an initializing create rolls back', async () => { + const mcpStarted = deferred(); + const mcpReady = deferred(); + const failure = new Error('MCP initialization failed'); + const svc = build([ + stubPair(IAgentLifecycleService, { + ...agentLifecycleStub(), + ensureMcpReady: () => { + mcpStarted.resolve(); + return mcpReady.promise; + }, + }), + ]); + + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await mcpStarted.promise; + const rejected = expect(create).rejects.toBe(failure); + const resume = svc.resume('s1'); + + mcpReady.reject(failure); + await rejected; + await expect(resume).resolves.toBeUndefined(); + expect(svc.get('s1')).toBeUndefined(); + }); + + it('removes fresh persisted state when create initialization fails', async () => { + const root = await makeTmpRoot(); + const storage = new FileStorageService(root); + const mcpStarted = deferred(); + const mcpReady = deferred(); + const failure = new Error('MCP initialization failed'); + let mcpAttempts = 0; + const queryStore = statefulQueryStore(); + registerScopedService( + LifecycleScope.Session, + ISessionMetadata, + SessionMetadata, + InstantiationType.Delayed, + 'sessionMetadata', + ); + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + stubPair(IFileSystemStorageService, storage), + [IAtomicDocumentStore, new SyncDescriptor(JsonAtomicDocumentStore)], + stubPair(IQueryStore, queryStore), + stubPair(IFlagService, stubFlag(false)), + stubPair(ILogService, stubLog()), + [ISessionIndex, new SyncDescriptor(FileSessionIndex)], + stubPair(IAgentLifecycleService, { + ...agentLifecycleWithMainStub(), + ensureMcpReady: () => { + mcpAttempts += 1; + if (mcpAttempts === 1) { + mcpStarted.resolve(); + return mcpReady.promise; + } + return Promise.resolve(); + }, + }), + ]); + const index = host!.app.accessor.get(ISessionIndex); + const sessionDir = join(root, 'sessions', 'wd_stub', 's1'); + + const creating = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await mcpStarted.promise; + await expect(index.get('s1')).resolves.toMatchObject({ id: 's1' }); + const rejected = expect(creating).rejects.toBe(failure); + mcpReady.reject(failure); + await rejected; + + await expect(stat(sessionDir)).rejects.toThrow(); + await expect(index.get('s1')).resolves.toBeUndefined(); + await expect(index.list({ includeArchived: true })).resolves.toMatchObject({ items: [] }); + await expect(svc.resume('s1')).resolves.toBeUndefined(); + + const retried = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + expect(retried.id).toBe('s1'); + }); + + it('waits for queued metadata mirroring before removing its owned directory', async () => { + const root = await makeTmpRoot(); + const storage = new FileStorageService(root); + const queryStore = statefulQueryStore(); + const putStarted = deferred(); + const releasePut = deferred(); + const idleStarted = deferred(); + const mcpStarted = deferred(); + const mcpReady = deferred(); + const failure = new Error('MCP initialization failed'); + const pendingMirror = (async () => { + putStarted.resolve(); + await releasePut.promise; + await queryStore.put('session', 's1', { + id: 's1', + workspaceId: 'wd_stub', + createdAt: 1, + updatedAt: 1, + archived: false, + }); + })(); + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + stubPair(IFileSystemStorageService, storage), + stubPair(IQueryStore, queryStore), + stubPair(IFlagService, stubFlag(true)), + stubPair(ILogService, stubLog()), + [ISessionIndex, new SyncDescriptor(FileSessionIndex)], + stubPair(ISessionMetadata, { + ...metadataStub(), + whenIdle: async () => { + idleStarted.resolve(); + await pendingMirror; + }, + }), + stubPair(IAgentLifecycleService, { + ...agentLifecycleWithMainStub(), + ensureMcpReady: () => { + mcpStarted.resolve(); + return mcpReady.promise; + }, + }), + ]); + const index = host!.app.accessor.get(ISessionIndex); + const sessionDir = join(root, 'sessions', 'wd_stub', 's1'); + let settled = false; - const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }).then(() => { + const creating = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }).finally(() => { settled = true; }); + await putStarted.promise; + await mcpStarted.promise; + mcpReady.reject(failure); + await idleStarted.promise; + await Promise.resolve(); + expect(settled).toBe(false); + await expect(stat(sessionDir)).resolves.toBeDefined(); + + releasePut.resolve(); + await expect(creating).rejects.toBe(failure); + await expect(stat(sessionDir)).rejects.toThrow(); + await expect(queryStore.get('session', 's1')).resolves.toMatchObject({ id: 's1' }); + await expect(index.get('s1')).resolves.toBeUndefined(); + }); - await tick(); + it('waits for append-log writes before removing its owned directory', async () => { + const root = await makeTmpRoot(); + const flushStarted = deferred(); + const releaseFlush = deferred(); + const failure = new Error('MCP initialization failed'); + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + stubPair(IAppendLogStore, { + ...appendLogStoreStub(), + flush: () => { + flushStarted.resolve(); + return releaseFlush.promise; + }, + }), + stubPair(IAgentLifecycleService, { + ...agentLifecycleWithMainStub(), + ensureMcpReady: () => Promise.reject(failure), + }), + ]); + const sessionDir = join(root, 'sessions', 'wd_stub', 's1'); + + let settled = false; + const creating = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }).finally(() => { + settled = true; + }); + await flushStarted.promise; + await Promise.resolve(); expect(settled).toBe(false); + await expect(stat(sessionDir)).resolves.toBeDefined(); - resolveMcpReady?.(); - await create; - expect(settled).toBe(true); + releaseFlush.resolve(); + await expect(creating).rejects.toBe(failure); + await expect(stat(sessionDir)).rejects.toThrow(); + }); + + it('reports cleanup failure without releasing an unremoved owned directory', async () => { + const root = await makeTmpRoot(); + const sessionDir = join(root, 'sessions', 'wd_stub', 's1'); + const realHostFs = new HostFileSystem(); + const cleanupFailure = new Error('session directory cleanup failed'); + const failingHostFs = new Proxy(realHostFs, { + get(target, property, receiver) { + if (property === 'remove') { + return (path: string): Promise => + path === sessionDir ? Promise.reject(cleanupFailure) : target.remove(path); + } + const value = Reflect.get(target, property, receiver) as unknown; + return typeof value === 'function' ? value.bind(target) : value; + }, + }) as IHostFileSystem; + const initializationFailure = new Error('MCP initialization failed'); + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + stubPair(IHostFileSystem, failingHostFs), + stubPair(IAgentLifecycleService, { + ...agentLifecycleWithMainStub(), + ensureMcpReady: () => Promise.reject(initializationFailure), + }), + ]); + + const failure = await svc + .create({ sessionId: 's1', workDir: '/tmp/proj' }) + .then(() => undefined, (error: unknown) => error); + + expect(failure).toBeInstanceOf(AggregateError); + expect((failure as AggregateError).errors).toEqual([ + initializationFailure, + cleanupFailure, + ]); + await expect(stat(sessionDir)).resolves.toBeDefined(); + await expect( + svc.create({ sessionId: 's1', workDir: '/tmp/proj' }), + ).rejects.toMatchObject({ code: ErrorCodes.SESSION_ALREADY_EXISTS }); + }); + + it('does not delete a session directory that wins the exclusive create race', async () => { + const root = await makeTmpRoot(); + const realHostFs = new HostFileSystem(); + const sessionDir = join(root, 'sessions', 'wd_stub', 's1'); + const sessionParent = join(root, 'sessions', 'wd_stub'); + const statePath = join(sessionDir, 'state.json'); + const externalState = JSON.stringify({ owner: 'external' }); + const queryStore = statefulQueryStore(); + const externalSummary = { id: 's1', workspaceId: 'wd_stub', cwd: '/tmp/proj' }; + await queryStore.put('session', 's1', externalSummary); + let injected = false; + const racingHostFs = new Proxy(realHostFs, { + get(target, property, receiver) { + if (property !== 'mkdir') return Reflect.get(target, property, receiver); + return async (path: string, options?: { readonly recursive?: boolean }): Promise => { + await target.mkdir(path, options); + if (!injected && path === sessionParent && options?.recursive === true) { + injected = true; + await target.mkdir(sessionDir); + await target.writeText(statePath, externalState); + } + }; + }, + }); + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + stubPair(IHostFileSystem, racingHostFs), + stubPair(IQueryStore, queryStore), + ]); + + await expect( + svc.create({ sessionId: 's1', workDir: '/tmp/proj' }), + ).rejects.toMatchObject({ code: ErrorCodes.SESSION_ALREADY_EXISTS }); + + await expect(readFile(statePath, 'utf8')).resolves.toBe(externalState); + await expect(queryStore.get('session', 's1')).resolves.toEqual(externalSummary); + }); + + it('rejects create without deleting an existing persisted session', async () => { + const root = await makeTmpRoot(); + const storage = new FileStorageService(root); + const queryStore = statefulQueryStore(); + registerScopedService( + LifecycleScope.Session, + ISessionMetadata, + SessionMetadata, + InstantiationType.Delayed, + 'sessionMetadata', + ); + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + stubPair(IFileSystemStorageService, storage), + [IAtomicDocumentStore, new SyncDescriptor(JsonAtomicDocumentStore)], + stubPair(IQueryStore, queryStore), + stubPair(IFlagService, stubFlag(true)), + stubPair(ILogService, stubLog()), + [ISessionIndex, new SyncDescriptor(FileSessionIndex)], + stubPair(IAgentLifecycleService, agentLifecycleWithMainStub()), + ]); + const index = host!.app.accessor.get(ISessionIndex); + const statePath = join(root, 'sessions', 'wd_stub', 's1', 'state.json'); + await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await svc.close('s1'); + const stateBefore = await readFile(statePath, 'utf8'); + await expect(index.get('s1')).resolves.toMatchObject({ id: 's1' }); + + await expect( + svc.create({ sessionId: 's1', workDir: '/tmp/proj' }), + ).rejects.toMatchObject({ code: ErrorCodes.SESSION_ALREADY_EXISTS }); + + await expect(readFile(statePath, 'utf8')).resolves.toBe(stateBefore); + await expect(queryStore.get('session', 's1')).resolves.toMatchObject({ id: 's1' }); + await expect(index.get('s1')).resolves.toMatchObject({ id: 's1' }); + }); + + it('preserves an existing persisted session when resume initialization fails', async () => { + const root = await makeTmpRoot(); + const storage = new FileStorageService(root); + const mcpStarted = deferred(); + const mcpReady = deferred(); + const failure = new Error('resume MCP initialization failed'); + let mcpAttempts = 0; + registerScopedService( + LifecycleScope.Session, + ISessionMetadata, + SessionMetadata, + InstantiationType.Delayed, + 'sessionMetadata', + ); + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + stubPair(IFileSystemStorageService, storage), + [IAtomicDocumentStore, new SyncDescriptor(JsonAtomicDocumentStore)], + stubPair(IQueryStore, statefulQueryStore()), + stubPair(IFlagService, stubFlag(false)), + stubPair(ILogService, stubLog()), + [ISessionIndex, new SyncDescriptor(FileSessionIndex)], + stubPair(IAgentLifecycleService, { + ...agentLifecycleWithMainStub(), + ensureMcpReady: () => { + mcpAttempts += 1; + if (mcpAttempts === 2) { + mcpStarted.resolve(); + return mcpReady.promise; + } + return Promise.resolve(); + }, + }), + ]); + const index = host!.app.accessor.get(ISessionIndex); + const statePath = join(root, 'sessions', 'wd_stub', 's1', 'state.json'); + await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await svc.close('s1'); + const stateBefore = await readFile(statePath, 'utf8'); + + const resuming = svc.resume('s1'); + await mcpStarted.promise; + const rejected = expect(resuming).rejects.toBe(failure); + mcpReady.reject(failure); + await rejected; + + await expect(readFile(statePath, 'utf8')).resolves.toBe(stateBefore); + await expect(index.get('s1')).resolves.toMatchObject({ id: 's1' }); + const retried = await svc.resume('s1'); + expect(retried?.id).toBe('s1'); + }); + + it('honors close after an in-flight session reaches publication', async () => { + const firstMcpStarted = deferred(); + const firstMcp = deferred(); + let attempts = 0; + const svc = build([ + stubPair(IAgentLifecycleService, { + ...agentLifecycleStub(), + ensureMcpReady: () => { + attempts += 1; + if (attempts === 1) { + firstMcpStarted.resolve(); + return firstMcp.promise; + } + return Promise.resolve(); + }, + }), + ]); + const closed: string[] = []; + svc.onDidCloseSession((event) => closed.push(event.sessionId)); + + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await firstMcpStarted.promise; + const close = svc.close('s1'); + const rejected = expect(create).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, + }); + firstMcp.resolve(); + await rejected; + await close; + expect(closed).toEqual(['s1']); + await expect( + svc.create({ sessionId: 's1', workDir: '/tmp/proj' }), + ).rejects.toMatchObject({ code: ErrorCodes.SESSION_ALREADY_EXISTS }); + }); + + it('honors close requested while workspace registration blocks handle creation', async () => { + const workspaceStarted = deferred(); + const workspaceReady = deferred(); + const registry = workspaceRegistryStub(); + const svc = build([ + stubPair(IWorkspaceRegistry, { + ...registry, + createOrTouch: () => { + workspaceStarted.resolve(); + return workspaceReady.promise; + }, + }), + ]); + const closed: string[] = []; + svc.onDidCloseSession((event) => closed.push(event.sessionId)); + + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await workspaceStarted.promise; + const close = svc.close('s1'); + const rejected = expect(create).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, + }); + expect(closed).toEqual([]); + expect(svc.get('s1')).toBeUndefined(); + workspaceReady.resolve({ + id: 'wd_stub', + root: '/tmp/proj', + name: 'stub', + createdAt: 0, + lastOpenedAt: 0, + }); + + await close; + await rejected; + expect(closed).toEqual(['s1']); + expect(svc.get('s1')).toBeUndefined(); + }); + + it('honors archive requested while the host environment blocks handle creation', async () => { + const localConfigRead = deferred(); + const hostReady = deferred(); + const archivedValues: boolean[] = []; + const localConfig = workspaceLocalConfigStub(); + const svc = build([ + stubPair(IHostEnvironment, { ...hostEnvironmentStub(), ready: hostReady.promise }), + stubPair(IWorkspaceLocalConfigService, { + ...localConfig, + resolveAdditionalDirs: (baseDir, dirs) => { + localConfigRead.resolve(); + return localConfig.resolveAdditionalDirs(baseDir, dirs); + }, + }), + stubPair(ISessionMetadata, { + ...metadataStub(), + setArchived: (value) => { + archivedValues.push(value); + return Promise.resolve(); + }, + }), + ]); + const archived: string[] = []; + svc.onDidArchiveSession((event) => archived.push(event.sessionId)); + + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await localConfigRead.promise; + const archive = svc.archive('s1'); + const rejected = expect(create).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, + }); + hostReady.resolve(); + + await archive; + await rejected; + expect(archivedValues).toEqual([true]); + expect(archived).toEqual(['s1']); + expect(svc.get('s1')).toBeUndefined(); + }); + + it('settles a pre-handle close when session initialization fails', async () => { + const workspaceStarted = deferred(); + const workspaceReady = deferred(); + const failure = new Error('workspace registration failed'); + const registry = workspaceRegistryStub(); + const svc = build([ + stubPair(IWorkspaceRegistry, { + ...registry, + createOrTouch: () => { + workspaceStarted.resolve(); + return workspaceReady.promise; + }, + }), + ]); + + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await workspaceStarted.promise; + const close = svc.close('s1'); + const rejected = expect(create).rejects.toBe(failure); + workspaceReady.reject(failure); + + await rejected; + await expect(close).resolves.toBeUndefined(); + expect(svc.get('s1')).toBeUndefined(); + }); + + it('allows only one same-id create initialization until failed cleanup releases the id', async () => { + const firstMcpStarted = deferred(); + const firstMcp = deferred(); + const failure = new Error('MCP initialization failed'); + let attempts = 0; + let workspaceTouches = 0; + const workspaceRegistry = workspaceRegistryStub(); + const svc = build([ + stubPair(IWorkspaceRegistry, { + ...workspaceRegistry, + createOrTouch: (root, name) => { + workspaceTouches += 1; + return workspaceRegistry.createOrTouch(root, name); + }, + }), + stubPair(IAgentLifecycleService, { + ...agentLifecycleStub(), + ensureMcpReady: () => { + attempts += 1; + if (attempts === 1) { + firstMcpStarted.resolve(); + return firstMcp.promise; + } + return Promise.resolve(); + }, + }), + ]); + + const firstCreate = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await firstMcpStarted.promise; + await expect(svc.create({ sessionId: 's1', workDir: '/tmp/proj' })).rejects.toMatchObject({ + code: ErrorCodes.SESSION_ALREADY_EXISTS, + }); + expect(attempts).toBe(1); + expect(workspaceTouches).toBe(1); + const rejected = expect(firstCreate).rejects.toBe(failure); + firstMcp.reject(failure); + await rejected; + + expect(svc.get('s1')).toBeUndefined(); + + const retried = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + expect(svc.get('s1')).toBe(retried); + expect(attempts).toBe(2); + expect(workspaceTouches).toBe(2); + }); + + it('keeps a resumed session published when its creation hook fails', async () => { + const failure = new Error('session creation hook failed'); + let attempts = 0; + const svc = build([ + stubPair(ISessionIndex, sessionIndexWithSummary('s1', '/tmp/proj')), + stubPair(IAgentLifecycleService, agentLifecycleWithMainStub()), + ]); + const hook = svc.hooks.onDidCreateSession.register('fail-once', async (_event, next) => { + attempts += 1; + if (attempts === 1) throw failure; + await next(); + }); + + await expect(svc.resume('s1')).rejects.toBe(failure); + const published = svc.get('s1'); + + expect(published).toBeDefined(); + await expect(svc.resume('s1')).resolves.toBe(published); + expect(attempts).toBe(1); + hook.dispose(); + }); + + it('allows a same-id create retry when failed-session agent cleanup rejects', async () => { + const flushStarted = deferred(); + const flush = deferred(); + const initializationFailure = new Error('session index flush failed'); + const removed: string[] = []; + let flushAttempts = 0; + const agentHandle = (id: string) => + ({ + id, + kind: LifecycleScope.Agent, + accessor: { get: () => ({}) }, + dispose: () => {}, + }) as unknown as IAgentScopeHandle; + const svc = build([ + stubPair(IAppendLogStore, { + ...appendLogStoreStub(), + flush: () => { + flushAttempts += 1; + if (flushAttempts === 1) { + flushStarted.resolve(); + return flush.promise; + } + return Promise.resolve(); + }, + }), + stubPair(IAgentLifecycleService, { + ...agentLifecycleStub(), + list: () => [agentHandle('main'), agentHandle('subagent')], + remove: (id: string) => { + removed.push(id); + return id === 'main' + ? Promise.reject(new Error('agent removal failed')) + : Promise.resolve(); + }, + }), + ]); + + const creation = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + await flushStarted.promise; + const rejected = expect(creation).rejects.toBe(initializationFailure); + flush.reject(initializationFailure); + await rejected; + + expect(removed).toHaveLength(2); + expect(removed).toEqual(expect.arrayContaining(['main', 'subagent'])); + expect(svc.get('s1')).toBeUndefined(); + + const retried = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + expect(svc.get('s1')).toBe(retried); }); it('hides a session from get/list until its resume finishes', async () => { + const mcpStarted = deferred(); let resolveMcpReady: (() => void) | undefined; const mcpReady = new Promise((resolve) => { resolveMcpReady = resolve; @@ -789,12 +1686,15 @@ describe('SessionLifecycleService', () => { stubPair(ISessionIndex, sessionIndexWithSummary('s1', '/tmp/proj')), stubPair(IAgentLifecycleService, { ...agentLifecycleWithMainStub(), - ensureMcpReady: () => mcpReady, + ensureMcpReady: () => { + mcpStarted.resolve(); + return mcpReady; + }, }), ]); const resumed = svc.resume('s1'); - await tick(); + await mcpStarted.promise; // materialize has registered the handle in `sessions` and is now blocked on // ensureMcpReady with `resuming` set — the handle must not be observable yet. @@ -1012,6 +1912,214 @@ describe('SessionLifecycleService', () => { }); } + it('allows a published source creation hook to fork the same session', async () => { + let forked: ISessionScopeHandle | undefined; + const svc = build([workspaceGetStub()]); + const hook = svc.hooks.onDidCreateSession.register( + 'fork-created-source', + async (event, next) => { + if (event.sessionId === 'src') { + forked = await svc.fork({ sourceSessionId: event.sessionId, newSessionId: 'dst' }); + } + await next(); + }, + ); + + await svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + + expect(forked).toBeDefined(); + expect(svc.get('dst')).toBe(forked); + hook.dispose(); + }); + + it('keeps a fork target published when its creation hook fails', async () => { + const failure = new Error('fork creation hook failed'); + let published: ISessionScopeHandle | undefined; + const svc = build([workspaceGetStub()]); + const hook = svc.hooks.onDidCreateSession.register( + 'fail-fork-session', + async (event, next) => { + if (event.source === 'fork') { + published = event.handle; + throw failure; + } + await next(); + }, + ); + await svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + + await expect( + svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }), + ).rejects.toBe(failure); + + expect(published).toBeDefined(); + expect(svc.get('dst')).toBe(published); + await expect(svc.resume('dst')).resolves.toBe(published); + hook.dispose(); + }); + + it('waits for a creating source to publish before forking it', async () => { + const mcpStarted = deferred(); + const mcpReady = deferred(); + let mcpAttempts = 0; + const indexGet = vi.fn(() => Promise.resolve(undefined)); + const svc = build([ + workspaceGetStub(), + stubPair(ISessionIndex, { ...sessionIndexStub(), get: indexGet }), + stubPair(IAgentLifecycleService, { + ...agentLifecycleStub(), + ensureMcpReady: () => { + mcpAttempts += 1; + if (mcpAttempts === 1) { + mcpStarted.resolve(); + return mcpReady.promise; + } + return Promise.resolve(); + }, + }), + ]); + + const creating = svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + await mcpStarted.promise; + const forked = svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + + expect(indexGet).not.toHaveBeenCalled(); + mcpReady.resolve(); + const source = await creating; + const target = await forked; + + expect(svc.get('src')).toBe(source); + expect(target.id).toBe('dst'); + }); + + it('waits for a resuming source to publish before forking it', async () => { + const mcpStarted = deferred(); + const mcpReady = deferred(); + let mcpAttempts = 0; + const sourceIndex = sessionIndexWithSummary('src', '/tmp/proj', 'wd_stub'); + const indexGet = vi.fn(sourceIndex.get.bind(sourceIndex)); + const svc = build([ + workspaceGetStub(), + stubPair(ISessionIndex, { ...sourceIndex, get: indexGet }), + stubPair(IAgentLifecycleService, { + ...agentLifecycleWithMainStub(), + ensureMcpReady: () => { + mcpAttempts += 1; + if (mcpAttempts === 1) { + mcpStarted.resolve(); + return mcpReady.promise; + } + return Promise.resolve(); + }, + }), + ]); + + const resuming = svc.resume('src'); + await mcpStarted.promise; + const indexCallsBeforeFork = indexGet.mock.calls.length; + const forked = svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + + expect(indexGet).toHaveBeenCalledTimes(indexCallsBeforeFork); + mcpReady.resolve(); + const source = await resuming; + const target = await forked; + + expect(svc.get('src')).toBe(source); + expect(target.id).toBe('dst'); + }); + + it('does not fork a source whose initialization fails', async () => { + const mcpStarted = deferred(); + const mcpReady = deferred(); + const failure = new Error('source MCP initialization failed'); + const indexGet = vi.fn(() => Promise.resolve(undefined)); + const svc = build([ + workspaceGetStub(), + stubPair(ISessionIndex, { ...sessionIndexStub(), get: indexGet }), + stubPair(IAgentLifecycleService, { + ...agentLifecycleStub(), + ensureMcpReady: () => { + mcpStarted.resolve(); + return mcpReady.promise; + }, + }), + ]); + + const creating = svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + await mcpStarted.promise; + const forked = svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + const createRejected = expect(creating).rejects.toBe(failure); + const forkRejected = expect(forked).rejects.toMatchObject({ + code: ErrorCodes.SESSION_NOT_FOUND, + }); + + expect(indexGet).not.toHaveBeenCalled(); + mcpReady.reject(failure); + await createRejected; + await forkRejected; + + expect(svc.get('src')).toBeUndefined(); + expect(svc.get('dst')).toBeUndefined(); + expect(svc.list()).toEqual([]); + }); + + it('waits for source close to settle before using the persisted fallback', async () => { + const createHookStarted = deferred(); + const finishCreateHook = deferred(); + const closeHookStarted = deferred(); + const finishCloseHook = deferred(); + const sourceIndex = sessionIndexWithSummary('src', '/tmp/proj', 'wd_stub'); + const indexGet = vi.fn(sourceIndex.get.bind(sourceIndex)); + const svc = build([ + workspaceGetStub(), + stubPair(ISessionIndex, { ...sourceIndex, get: indexGet }), + ]); + const createHook = svc.hooks.onDidCreateSession.register( + 'hold-created-source', + async (event, next) => { + if (event.sessionId === 'src') { + createHookStarted.resolve(); + await finishCreateHook.promise; + } + await next(); + }, + ); + const closeHook = svc.hooks.onWillCloseSession.register( + 'hold-closing-source', + async (event, next) => { + if (event.sessionId === 'src') { + closeHookStarted.resolve(); + await finishCloseHook.promise; + } + await next(); + }, + ); + + const creating = svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + const createRejected = expect(creating).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, + }); + await createHookStarted.promise; + const closing = svc.close('src'); + await closeHookStarted.promise; + const indexCallsBeforeFork = indexGet.mock.calls.length; + const forked = svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + + expect(indexGet).toHaveBeenCalledTimes(indexCallsBeforeFork); + finishCreateHook.resolve(); + await createRejected; + expect(indexGet).toHaveBeenCalledTimes(indexCallsBeforeFork); + + finishCloseHook.resolve(); + await closing; + const target = await forked; + + expect(target.id).toBe('dst'); + expect(svc.get('src')).toBeUndefined(); + createHook.dispose(); + closeHook.dispose(); + }); + it('copies blobs, plans, background tasks, and media originals into the fork', async () => { const root = await makeTmpRoot(); const svc = build([ @@ -1097,6 +2205,239 @@ describe('SessionLifecycleService', () => { ); }); + it('removes a fork target when materialization fails after metadata is persisted', async () => { + const root = await makeTmpRoot(); + const bootstrap = tmpBootstrapStub(root); + const storage = new FileStorageService(root); + const queryStore = stubQueryStore(); + const flags = stubFlag(false); + const log = stubLog(); + const targetMcpStarted = deferred(); + const targetMcp = deferred(); + const failure = new Error('target MCP initialization failed'); + let mcpAttempts = 0; + registerScopedService( + LifecycleScope.Session, + ISessionMetadata, + SessionMetadata, + InstantiationType.Delayed, + 'sessionMetadata', + ); + const svc = build([ + stubPair(IBootstrapService, bootstrap), + workspaceGetStub(), + stubPair(IFileSystemStorageService, storage), + [IAtomicDocumentStore, new SyncDescriptor(JsonAtomicDocumentStore)], + stubPair(IQueryStore, queryStore), + stubPair(IFlagService, flags), + stubPair(ILogService, log), + [ISessionIndex, new SyncDescriptor(FileSessionIndex)], + stubPair(IAgentLifecycleService, { + ...agentLifecycleStub(), + ensureMcpReady: () => { + mcpAttempts += 1; + if (mcpAttempts === 2) { + targetMcpStarted.resolve(); + return targetMcp.promise; + } + return Promise.resolve(); + }, + }), + ]); + const index = host!.app.accessor.get(ISessionIndex); + await svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + + const forked = svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + await targetMcpStarted.promise; + await expect(index.get('dst')).resolves.toMatchObject({ id: 'dst' }); + const rejected = expect(forked).rejects.toBe(failure); + targetMcp.reject(failure); + await rejected; + + await expect(stat(join(root, 'sessions', 'wd_stub', 'dst'))).rejects.toThrow(); + await expect(index.get('dst')).resolves.toBeUndefined(); + + const retried = await svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + expect(retried.id).toBe('dst'); + }); + + it('keeps a published fork resumable when close waits during initialization', async () => { + const root = await makeTmpRoot(); + const bootstrap = tmpBootstrapStub(root); + const storage = new FileStorageService(root); + const targetMcpStarted = deferred(); + const targetMcp = deferred(); + let mcpAttempts = 0; + registerScopedService( + LifecycleScope.Session, + ISessionMetadata, + SessionMetadata, + InstantiationType.Delayed, + 'sessionMetadata', + ); + const svc = build([ + stubPair(IBootstrapService, bootstrap), + workspaceGetStub(), + stubPair(IFileSystemStorageService, storage), + [IAtomicDocumentStore, new SyncDescriptor(JsonAtomicDocumentStore)], + stubPair(IQueryStore, stubQueryStore()), + stubPair(IFlagService, stubFlag(false)), + stubPair(ILogService, stubLog()), + [ISessionIndex, new SyncDescriptor(FileSessionIndex)], + stubPair(IAgentLifecycleService, { + ...agentLifecycleWithMainStub(), + ensureMcpReady: () => { + mcpAttempts += 1; + if (mcpAttempts === 2) { + targetMcpStarted.resolve(); + return targetMcp.promise; + } + return Promise.resolve(); + }, + }), + ]); + const index = host!.app.accessor.get(ISessionIndex); + await svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + + const forked = svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + await targetMcpStarted.promise; + const close = svc.close('dst'); + const forkRejected = expect(forked).rejects.toMatchObject({ + code: ErrorCodes.SESSION_CLOSED, + }); + targetMcp.resolve(); + await forkRejected; + await close; + + await expect(index.get('dst')).resolves.toMatchObject({ id: 'dst' }); + await expect( + stat(join(root, 'sessions', 'wd_stub', 'dst', 'state.json')), + ).resolves.toBeDefined(); + const resumed = await svc.resume('dst'); + expect(resumed?.id).toBe('dst'); + }); + + it('preserves a same-id session that registers during the fork collision check', async () => { + const collisionCheckStarted = deferred(); + const releaseCollisionCheck = deferred(); + const index: ISessionIndex = { + ...sessionIndexStub(), + get: async (sessionId: string) => { + if (sessionId === 'dst') { + collisionCheckStarted.resolve(); + await releaseCollisionCheck.promise; + } + return undefined; + }, + }; + const svc = build([workspaceGetStub(), stubPair(ISessionIndex, index)]); + await svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + + const forked = svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + await collisionCheckStarted.promise; + const replacement = await svc.create({ sessionId: 'dst', workDir: '/tmp/proj' }); + const rejected = expect(forked).rejects.toMatchObject({ + code: ErrorCodes.SESSION_ALREADY_EXISTS, + }); + releaseCollisionCheck.resolve(); + await rejected; + + expect(svc.get('dst')).toBe(replacement); + expect(replacement.accessor.get(ISessionContext).sessionId).toBe('dst'); + }); + + it('keeps the fork destination exclusively owned through failed copy rollback', async () => { + const root = await makeTmpRoot(); + const srcDir = join(root, 'sessions', 'wd_stub', 'src'); + const dstDir = join(root, 'sessions', 'wd_stub', 'dst'); + const storage = new FileStorageService(root); + const copyStarted = deferred(); + const copy = deferred(); + const rollbackStarted = deferred(); + const releaseRollback = deferred(); + const copyFailure = new Error('fork wire copy failed'); + const realHostFs = new HostFileSystem(); + const gatedHostFs = new Proxy(realHostFs, { + get(target, property, receiver) { + if (property === 'remove') { + return async (path: string): Promise => { + if (path === dstDir) { + rollbackStarted.resolve(); + await releaseRollback.promise; + } + await target.remove(path); + }; + } + const value = Reflect.get(target, property, receiver) as unknown; + return typeof value === 'function' ? value.bind(target) : value; + }, + }) as IHostFileSystem; + registerScopedService( + LifecycleScope.Session, + ISessionMetadata, + SessionMetadata, + InstantiationType.Delayed, + 'sessionMetadata', + ); + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + workspaceGetStub(), + stubPair(IFileSystemStorageService, storage), + [IAtomicDocumentStore, new SyncDescriptor(JsonAtomicDocumentStore)], + stubPair(IQueryStore, stubQueryStore()), + stubPair(IFlagService, stubFlag(false)), + stubPair(ILogService, stubLog()), + [ISessionIndex, new SyncDescriptor(FileSessionIndex)], + stubPair(IHostFileSystem, gatedHostFs), + stubPair(IAppendLogStore, { + ...appendLogStoreStub(), + rewrite: () => { + copyStarted.resolve(); + return copy.promise; + }, + }), + ]); + const source = await svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + await source.accessor.get(ISessionMetadata).registerAgent('main', { + homedir: join(srcDir, 'agents', 'main'), + type: 'main', + }); + await mkdir(join(srcDir, 'agents', 'main', 'plans'), { recursive: true }); + await writeFile(join(srcDir, 'agents', 'main', 'plans', 'p1.md'), '# plan'); + + const forked = svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }); + await copyStarted.promise; + const targetStateBefore = await readFile(join(dstDir, 'state.json'), 'utf8'); + await expect(svc.create({ sessionId: 'dst', workDir: '/tmp/proj' })).rejects.toMatchObject({ + code: ErrorCodes.SESSION_ALREADY_EXISTS, + }); + await expect(readFile(join(dstDir, 'state.json'), 'utf8')).resolves.toBe(targetStateBefore); + await expect(readFile(join(dstDir, 'agents', 'main', 'plans', 'p1.md'), 'utf8')).resolves.toBe( + '# plan', + ); + + const forkRejected = expect(forked).rejects.toBe(copyFailure); + copy.reject(copyFailure); + await rollbackStarted.promise; + await expect(svc.create({ sessionId: 'dst', workDir: '/tmp/proj' })).rejects.toMatchObject({ + code: ErrorCodes.SESSION_ALREADY_EXISTS, + }); + + releaseRollback.resolve(); + await forkRejected; + await expect(stat(dstDir)).rejects.toThrow(); + const index = host!.app.accessor.get(ISessionIndex); + await expect(index.get('dst')).resolves.toBeUndefined(); + + const replacement = await svc.create({ sessionId: 'dst', workDir: '/tmp/proj' }); + const replacementState = JSON.parse( + await readFile(join(dstDir, 'state.json'), 'utf8'), + ) as Record; + expect(svc.get('dst')).toBe(replacement); + expect(replacementState).not.toHaveProperty('forkedFrom'); + await expect(stat(join(dstDir, 'agents', 'main', 'plans', 'p1.md'))).rejects.toThrow(); + }); + it('duplicates the source session cron tasks for the fork', async () => { const root = await makeTmpRoot(); const cron = cronStoreStub([ diff --git a/packages/agent-core-v2/test/session/sessionMetadata/sessionMetadata.test.ts b/packages/agent-core-v2/test/session/sessionMetadata/sessionMetadata.test.ts index eb590436bc..b103ed9cea 100644 --- a/packages/agent-core-v2/test/session/sessionMetadata/sessionMetadata.test.ts +++ b/packages/agent-core-v2/test/session/sessionMetadata/sessionMetadata.test.ts @@ -115,4 +115,74 @@ describe('SessionMetadata', () => { 'agent-1', ]); }); + + it('whenIdle waits for a queued read-model write', async () => { + let startPut!: () => void; + const putStarted = new Promise((resolve) => { + startPut = resolve; + }); + let releasePut!: () => void; + const putReleased = new Promise((resolve) => { + releasePut = resolve; + }); + ix.stub(IFlagService, stubFlag(true)); + ix.stub(IQueryStore, { + ...stubQueryStore(), + put: async () => { + startPut(); + await putReleased; + }, + }); + const meta = ix.get(ISessionMetadata); + + const updating = meta.update({ title: 'queued' }); + await putStarted; + let idle = false; + const waiting = meta.whenIdle().then(() => { + idle = true; + }); + await Promise.resolve(); + expect(idle).toBe(false); + + releasePut(); + await updating; + await waiting; + expect(idle).toBe(true); + }); + + it('whenIdle waits for initial metadata persistence', async () => { + let startWrite!: () => void; + const writeStarted = new Promise((resolve) => { + startWrite = resolve; + }); + let releaseWrite!: () => void; + const writeReleased = new Promise((resolve) => { + releaseWrite = resolve; + }); + ix.stub(IAtomicDocumentStore, { + _serviceBrand: undefined, + get: async () => undefined, + set: async () => { + startWrite(); + await writeReleased; + }, + delete: async () => {}, + list: async () => [], + watch: () => () => ({ dispose: () => {} }), + acquire: () => ({ dispose: () => {} }), + }); + const meta = ix.get(ISessionMetadata); + + let idle = false; + const waiting = meta.whenIdle().then(() => { + idle = true; + }); + await writeStarted; + await Promise.resolve(); + expect(idle).toBe(false); + + releaseWrite(); + await waiting; + expect(idle).toBe(true); + }); }); diff --git a/packages/agent-core-v2/test/tool/tool.test.ts b/packages/agent-core-v2/test/tool/tool.test.ts index fa73bb7e84..0cddeb1290 100644 --- a/packages/agent-core-v2/test/tool/tool.test.ts +++ b/packages/agent-core-v2/test/tool/tool.test.ts @@ -334,6 +334,7 @@ function sessionMetadataStub(agents: Readonly>): ISess _serviceBrand: undefined, ready: Promise.resolve(), onDidChangeMetadata: Event.None as ISessionMetadata['onDidChangeMetadata'], + whenIdle: async () => {}, read: async () => ({ id: 'test-session', createdAt: 0,