From 256fa4d2d75cab9162f5fc203fe67151c8ef62bd Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 14 Jul 2026 05:13:25 +0800 Subject: [PATCH 1/8] fix(agent-core-v2): roll back failed session initialization --- .changeset/clean-session-init-failures.md | 5 + .../sessionLifecycleService.ts | 188 +++++++++----- .../sessionLifecycle/sessionLifecycle.test.ts | 230 ++++++++++++++++++ 3 files changed, 364 insertions(+), 59 deletions(-) create mode 100644 .changeset/clean-session-init-failures.md diff --git a/.changeset/clean-session-init-failures.md b/.changeset/clean-session-init-failures.md new file mode 100644 index 0000000000..00a4576d05 --- /dev/null +++ b/.changeset/clean-session-init-failures.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/agent-core-v2": patch +--- + +Release partially initialized session scopes after creation or resume fails so the same session ID can be retried safely. diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts index dcc9cd14c3..3e2a410b92 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -13,6 +13,9 @@ * roots are remembered through `workspaceRegistry`. On create / fork the * session is also appended to the shared `session_index.jsonl` so v1 clients * (TUI, export) can discover sessions created by the v2 engine. + * Failed materialization, creation, resume, and fork attempts release only the + * Session handle they created. Fork rollback keeps the target id reserved + * until its partial directory has been removed. */ import { randomUUID } from 'node:crypto'; @@ -94,6 +97,7 @@ type MaterializeSessionOptions = Omit & { export class SessionLifecycleService extends Disposable implements ISessionLifecycleService { declare readonly _serviceBrand: undefined; private readonly sessions = new Map(); + private readonly forkRollbacks = new Map>(); private readonly _onDidCreateSession = this._register(new Emitter()); readonly onDidCreateSession: Event = this._onDidCreateSession.event; private readonly _onDidCloseSession = this._register(new Emitter()); @@ -135,14 +139,22 @@ 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(); + let handle: ISessionScopeHandle | undefined; + try { + 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(); + } + await this.announceCreated({ sessionId, handle, source: 'startup' }); + return handle; + } catch (error) { + if (handle !== undefined) { + await this.disposeFailedSession(sessionId, handle); + } + throw error; } - await this.announceCreated({ sessionId, handle, source: 'startup' }); - return handle; } private async materializeSession(opts: MaterializeSessionOptions): Promise { @@ -188,21 +200,21 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec 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); + try { + handle.accessor.get(ISessionActivityKernel); + if (additionalDirs.length > 0) { + handle.accessor.get(ISessionWorkspaceContext).setAdditionalDirs(additionalDirs); + } + await this.registerSessionHandle(opts.sessionId, handle); + await handle.accessor.get(ISessionMetadata).ready; + void handle.accessor.get(ISessionSkillCatalog).ready; + await handle.accessor.get(IAgentLifecycleService).ensureMcpReady(); + handle.accessor.get(ISessionExternalHooksService); + return handle; + } catch (error) { + await this.disposeFailedSession(opts.sessionId, handle); + throw error; } - this.sessions.set(opts.sessionId, handle); - await handle.accessor.get(ISessionMetadata).ready; - void handle.accessor.get(ISessionSkillCatalog).ready; - await handle.accessor.get(IAgentLifecycleService).ensureMcpReady(); - handle.accessor.get(ISessionExternalHooksService); - return handle; } private async appendSessionIndexEntry(sessionId: string, workDir: string): Promise { @@ -253,10 +265,14 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec 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 { @@ -272,25 +288,30 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec const workDir = summary.cwd ?? workspace?.root; if (workDir === undefined) return undefined; - 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); + let handle: ISessionScopeHandle | undefined; + try { + 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); + 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); + } + await this.announceCreated({ sessionId, handle, source: 'resume' }); + return handle; + } catch (error) { + if (handle !== undefined) { + await this.disposeFailedSession(sessionId, handle); + } + throw error; } - await this.announceCreated({ sessionId, handle, source: 'resume' }); - return handle; } list(): readonly ISessionScopeHandle[] { @@ -349,6 +370,69 @@ 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 { + handle.dispose(); + } catch {} + } + + private async registerSessionHandle( + sessionId: string, + handle: ISessionScopeHandle, + ): Promise { + while (true) { + const rollback = this.forkRollbacks.get(sessionId); + if (rollback === undefined) { + this.sessions.set(sessionId, handle); + return; + } + await rollback; + } + } + + private async rollbackFailedFork( + sessionId: string, + handle: ISessionScopeHandle, + sessionDir: string | undefined, + ): Promise { + if (this.sessions.get(sessionId) !== handle || sessionDir === undefined) { + await this.disposeFailedSession(sessionId, handle); + return; + } + let release!: () => void; + const rollback = new Promise((resolve) => { + release = resolve; + }); + this.forkRollbacks.set(sessionId, rollback); + this.sessions.delete(sessionId); + try { + await this.disposeFailedSession(sessionId, handle); + await this.hostFs.remove(sessionDir).catch(() => {}); + } finally { + if (this.forkRollbacks.get(sessionId) === rollback) { + this.forkRollbacks.delete(sessionId); + } + release(); + } + } + async fork(opts: ForkSessionOptions): Promise { const sourceId = opts.sourceSessionId; @@ -472,22 +556,8 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await this.announceCreated({ sessionId: targetId, handle: target, source: 'fork' }); 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) { + await this.rollbackFailedFork(targetId, target, targetSessionDir); } throw error; } finally { 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 8eb745c793..fe9ff2a9e9 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'; @@ -358,6 +368,20 @@ function agentLifecycleCapturingPlanSpy(opts: { mainPreexists?: boolean } = {}): return { lifecycle, enter, create }; } +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 }; +} + function tick(): Promise { return new Promise((resolve) => setTimeout(resolve, 0)); } @@ -770,6 +794,145 @@ describe('SessionLifecycleService', () => { expect(settled).toBe(true); }); + it('allows a same-id create retry after MCP initialization fails', async () => { + const firstMcpStarted = deferred(); + const firstMcp = deferred(); + const failure = new Error('MCP initialization failed'); + let attempts = 0; + const svc = build([ + 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; + const failedHandle = svc.get('s1'); + const rejected = expect(firstCreate).rejects.toBe(failure); + firstMcp.reject(failure); + await rejected; + + expect(svc.get('s1')).toBeUndefined(); + expect(() => failedHandle?.accessor.get(ISessionContext)).toThrow(/disposed/); + + const retried = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + expect(svc.get('s1')).toBe(retried); + }); + + it('allows a same-id resume retry after 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); + expect(svc.get('s1')).toBeUndefined(); + + const retried = await svc.resume('s1'); + expect(svc.get('s1')).toBe(retried); + hook.dispose(); + }); + + it('preserves a concurrent replacement when an older materialization fails', async () => { + const firstMcpStarted = deferred(); + const firstMcp = deferred(); + const failure = new Error('older MCP initialization failed'); + let attempts = 0; + const svc = build([ + 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; + const failedHandle = svc.get('s1'); + const replacement = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + const rejected = expect(firstCreate).rejects.toBe(failure); + firstMcp.reject(failure); + await rejected; + + expect(svc.get('s1')).toBe(replacement); + expect(svc.list()).toEqual([replacement]); + expect(() => failedHandle?.accessor.get(ISessionContext)).toThrow(/disposed/); + }); + + 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 failedHandle = svc.get('s1'); + 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(); + expect(() => failedHandle?.accessor.get(ISessionContext)).toThrow(/disposed/); + + 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 () => { let resolveMcpReady: (() => void) | undefined; const mcpReady = new Promise((resolve) => { @@ -1087,6 +1250,73 @@ describe('SessionLifecycleService', () => { ); }); + it('delays a same-id replacement until failed fork directory rollback finishes', async () => { + const root = await makeTmpRoot(); + const srcDir = join(root, 'sessions', 'wd_stub', 'src'); + const cleanupStarted = deferred(); + const releaseCleanup = deferred(); + const replacementReachedRegistration = deferred(); + const cleanupAgent = { + id: 'main', + kind: LifecycleScope.Agent, + accessor: { get: () => ({}) }, + dispose: () => {}, + } as unknown as IAgentScopeHandle; + const svc = build([ + stubPair(IBootstrapService, tmpBootstrapStub(root)), + workspaceGetStub(), + stubPair(ISessionMetadata, { + ...metadataStub(), + read: () => + Promise.resolve({ + agents: { main: { homedir: join(srcDir, 'agents', 'main') } }, + } as never), + }), + stubPair(IAgentLifecycleService, { + ...agentLifecycleStub(), + list: () => [cleanupAgent], + remove: () => { + cleanupStarted.resolve(); + return releaseCleanup.promise; + }, + }), + stubPair(ISessionWorkspaceContext, { + _serviceBrand: undefined, + workDir: '/tmp/proj', + additionalDirs: [], + setWorkDir: () => {}, + setAdditionalDirs: () => { + replacementReachedRegistration.resolve(); + }, + resolve: (path: string) => path, + isWithin: () => true, + assertAllowed: (path: string) => path, + addAdditionalDir: () => {}, + removeAdditionalDir: () => {}, + }), + ]); + await svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + 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 cleanupStarted.promise; + const forkRejected = expect(forked).rejects.toThrow('not implemented'); + const replacement = svc.create({ + sessionId: 'dst', + workDir: '/tmp/proj', + additionalDirs: ['/tmp/extra'], + }); + await replacementReachedRegistration.promise; + + expect(svc.get('dst')).toBeUndefined(); + + releaseCleanup.resolve(); + await forkRejected; + const replacementHandle = await replacement; + expect(svc.get('dst')).toBe(replacementHandle); + }); + it('duplicates the source session cron tasks for the fork', async () => { const root = await makeTmpRoot(); const cron = cronStoreStub([ From 7a5dd5318da9a54a2457e3e07b05969dc538fbf1 Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 14 Jul 2026 05:35:53 +0800 Subject: [PATCH 2/8] fix(agent-core-v2): roll back failed fork materialization --- .../sessionLifecycleService.ts | 17 ++- .../sessionLifecycle/sessionLifecycle.test.ts | 101 +++++++++++++++++- 2 files changed, 115 insertions(+), 3 deletions(-) diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts index 3e2a410b92..281e61bb76 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -92,6 +92,7 @@ import { type MaterializeSessionOptions = Omit & { readonly sessionId: string; readonly workspaceId?: string; + readonly forkDestination?: boolean; }; export class SessionLifecycleService extends Disposable implements ISessionLifecycleService { @@ -205,14 +206,18 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec if (additionalDirs.length > 0) { handle.accessor.get(ISessionWorkspaceContext).setAdditionalDirs(additionalDirs); } - await this.registerSessionHandle(opts.sessionId, handle); + await this.registerSessionHandle(opts.sessionId, handle, opts.forkDestination === true); await handle.accessor.get(ISessionMetadata).ready; void handle.accessor.get(ISessionSkillCatalog).ready; await handle.accessor.get(IAgentLifecycleService).ensureMcpReady(); handle.accessor.get(ISessionExternalHooksService); return handle; } catch (error) { - await this.disposeFailedSession(opts.sessionId, handle); + if (opts.forkDestination === true) { + await this.rollbackFailedFork(opts.sessionId, handle, ctx.sessionDir); + } else { + await this.disposeFailedSession(opts.sessionId, handle); + } throw error; } } @@ -396,10 +401,17 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec private async registerSessionHandle( sessionId: string, handle: ISessionScopeHandle, + exclusive: boolean, ): Promise { while (true) { const rollback = this.forkRollbacks.get(sessionId); if (rollback === undefined) { + if (exclusive && this.sessions.has(sessionId)) { + throw new Error2( + ErrorCodes.SESSION_ALREADY_EXISTS, + `Session "${sessionId}" already exists`, + ); + } this.sessions.set(sessionId, handle); return; } @@ -485,6 +497,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec target = await this.materializeSession({ sessionId: targetId, workDir: workspace.root, + forkDestination: true, }); const targetCtx = target.accessor.get(ISessionContext); targetSessionDir = targetCtx.sessionDir; 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 fe9ff2a9e9..17c975a7bc 100644 --- a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts +++ b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts @@ -15,10 +15,13 @@ 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, LifecycleScope, + type ScopeSeed, _clearScopedRegistryForTests, registerScopedService, } from '#/_base/di/scope'; @@ -26,6 +29,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'; @@ -50,8 +54,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'; @@ -60,8 +69,12 @@ 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 { @@ -78,6 +91,7 @@ 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) => @@ -464,7 +478,7 @@ 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(ISessionMetadata, metadataStub()), @@ -1250,6 +1264,91 @@ 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('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('delays a same-id replacement until failed fork directory rollback finishes', async () => { const root = await makeTmpRoot(); const srcDir = join(root, 'sessions', 'wd_stub', 'src'); From 833858caf78d53764ca5fc416e838b8508d0468d Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 14 Jul 2026 06:24:19 +0800 Subject: [PATCH 3/8] fix(agent-core-v2): claim session initialization ownership --- .changeset/clean-session-init-failures.md | 3 +- .../app/sessionLifecycle/sessionLifecycle.ts | 12 +- .../sessionLifecycleService.ts | 206 +++++++------ .../sessionLifecycle/sessionLifecycle.test.ts | 272 +++++++++++------- 4 files changed, 302 insertions(+), 191 deletions(-) diff --git a/.changeset/clean-session-init-failures.md b/.changeset/clean-session-init-failures.md index 00a4576d05..2a8903a095 100644 --- a/.changeset/clean-session-init-failures.md +++ b/.changeset/clean-session-init-failures.md @@ -1,5 +1,6 @@ --- "@moonshot-ai/agent-core-v2": patch +"@moonshot-ai/kimi-code": patch --- -Release partially initialized session scopes after creation or resume fails so the same session ID can be retried safely. +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/sessionLifecycle/sessionLifecycle.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts index a62bc45844..68c48eb887 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts @@ -99,17 +99,17 @@ 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 + * A session whose create, fork, or cold {@link resume} initialization is + * still in flight is intentionally NOT returned. Its metadata, agents, or + * replay may be incomplete. 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. + * invariant: every read path (`get` / {@link list} / {@link resume}) agrees + * an initializing session is not yet observable. */ 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}. + * initializing for the same reason as {@link get}. */ list(): readonly ISessionScopeHandle[]; /** diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts index 281e61bb76..2b5b25ec4b 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -14,8 +14,8 @@ * session is also appended to the shared `session_index.jsonl` so v1 clients * (TUI, export) can discover sessions created by the v2 engine. * Failed materialization, creation, resume, and fork attempts release only the - * Session handle they created. Fork rollback keeps the target id reserved - * until its partial directory has been removed. + * Session handle they created. Per-session initialization claims prevent + * overlapping create, resume, and fork attempts from sharing persistence. */ import { randomUUID } from 'node:crypto'; @@ -92,13 +92,19 @@ import { type MaterializeSessionOptions = Omit & { readonly sessionId: string; readonly workspaceId?: string; + readonly claim: SessionInitializationClaim; readonly forkDestination?: boolean; }; +interface SessionInitializationClaim { + readonly settled: Promise; + release(): void; +}; + export class SessionLifecycleService extends Disposable implements ISessionLifecycleService { declare readonly _serviceBrand: undefined; private readonly sessions = new Map(); - private readonly forkRollbacks = 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()); @@ -140,9 +146,10 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec async create(opts: CreateSessionOptions): Promise { const sessionId = opts.sessionId ?? createSessionId(); + const claim = this.claimInitialization(sessionId); let handle: ISessionScopeHandle | undefined; try { - handle = await this.materializeSession({ ...opts, sessionId }); + handle = await this.materializeSession({ ...opts, sessionId, claim }); await this.appendSessionIndexEntry(sessionId, opts.workDir); if (this.config.get(DEFAULT_PLAN_MODE_SECTION) === true) { const main = await ensureMainAgent(handle); @@ -155,6 +162,8 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await this.disposeFailedSession(sessionId, handle); } throw error; + } finally { + claim.release(); } } @@ -206,7 +215,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec if (additionalDirs.length > 0) { handle.accessor.get(ISessionWorkspaceContext).setAdditionalDirs(additionalDirs); } - await this.registerSessionHandle(opts.sessionId, handle, opts.forkDestination === true); + 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(); @@ -214,7 +223,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec return handle; } catch (error) { if (opts.forkDestination === true) { - await this.rollbackFailedFork(opts.sessionId, handle, ctx.sessionDir); + await this.rollbackFailedFork(opts.sessionId, handle, ctx.sessionDir, opts.claim); } else { await this.disposeFailedSession(opts.sessionId, handle); } @@ -244,12 +253,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } 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; + if (this.initializations.has(sessionId) || this.resuming.has(sessionId)) return undefined; return this.sessions.get(sessionId); } @@ -262,8 +266,10 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec // to complete. const inflight = this.resuming.get(sessionId); if (inflight !== undefined) return inflight; - const live = this.sessions.get(sessionId); - if (live !== undefined) return Promise.resolve(live); + if (!this.initializations.has(sessionId)) { + const live = this.sessions.get(sessionId); + if (live !== undefined) return Promise.resolve(live); + } const promise = this.doResume(sessionId) .catch((error: unknown) => { this.telemetry.track2('session_load_failed', { @@ -281,41 +287,60 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } 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; - - 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; + while (true) { + const initialization = this.initializations.get(sessionId); + if (initialization !== undefined) { + await initialization.settled; + continue; + } - let handle: ISessionScopeHandle | undefined; - try { - 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); - 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 live = this.sessions.get(sessionId); + if (live !== undefined) return live; + + 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) { + await racedInitialization.settled; + continue; } - await this.announceCreated({ sessionId, handle, source: 'resume' }); - return handle; - } catch (error) { - if (handle !== undefined) { - await this.disposeFailedSession(sessionId, handle); + const racedLive = this.sessions.get(sessionId); + if (racedLive !== undefined) return racedLive; + if (summary === undefined || workDir === undefined) return undefined; + + const claim = this.claimInitialization(sessionId); + 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); + } + await this.announceCreated({ sessionId, handle, source: 'resume' }); + return handle; + } catch (error) { + if (handle !== undefined) { + await this.disposeFailedSession(sessionId, handle); + } + throw error; + } finally { + claim.release(); } - throw error; } } @@ -324,7 +349,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec // 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); + if (!this.initializations.has(id) && !this.resuming.has(id)) ready.push(handle); } return ready; } @@ -398,51 +423,56 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } catch {} } - private async registerSessionHandle( + private claimInitialization(sessionId: string): SessionInitializationClaim { + if (this.sessions.has(sessionId) || this.initializations.has(sessionId)) { + throw sessionAlreadyExistsError(sessionId); + } + let released = false; + let settle!: () => void; + const settled = new Promise((resolve) => { + settle = resolve; + }); + const claim: SessionInitializationClaim = { + settled, + release: () => { + if (released) return; + released = true; + if (this.initializations.get(sessionId) === claim) { + this.initializations.delete(sessionId); + } + settle(); + }, + }; + this.initializations.set(sessionId, claim); + return claim; + } + + private registerSessionHandle( sessionId: string, handle: ISessionScopeHandle, - exclusive: boolean, - ): Promise { - while (true) { - const rollback = this.forkRollbacks.get(sessionId); - if (rollback === undefined) { - if (exclusive && this.sessions.has(sessionId)) { - throw new Error2( - ErrorCodes.SESSION_ALREADY_EXISTS, - `Session "${sessionId}" already exists`, - ); - } - this.sessions.set(sessionId, handle); - return; - } - await rollback; + claim: SessionInitializationClaim, + ): void { + if (this.initializations.get(sessionId) !== claim || this.sessions.has(sessionId)) { + throw sessionAlreadyExistsError(sessionId); } + this.sessions.set(sessionId, handle); } private async rollbackFailedFork( sessionId: string, handle: ISessionScopeHandle, sessionDir: string | undefined, + claim: SessionInitializationClaim, ): Promise { - if (this.sessions.get(sessionId) !== handle || sessionDir === undefined) { + if (this.initializations.get(sessionId) !== claim || sessionDir === undefined) { await this.disposeFailedSession(sessionId, handle); return; } - let release!: () => void; - const rollback = new Promise((resolve) => { - release = resolve; - }); - this.forkRollbacks.set(sessionId, rollback); - this.sessions.delete(sessionId); - try { - await this.disposeFailedSession(sessionId, handle); - await this.hostFs.remove(sessionDir).catch(() => {}); - } finally { - if (this.forkRollbacks.get(sessionId) === rollback) { - this.forkRollbacks.delete(sessionId); - } - release(); + if (this.sessions.get(sessionId) === handle) { + this.sessions.delete(sessionId); } + await this.disposeFailedSession(sessionId, handle); + await this.hostFs.remove(sessionDir).catch(() => {}); } async fork(opts: ForkSessionOptions): Promise { @@ -471,6 +501,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); @@ -487,16 +518,15 @@ 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); @@ -569,11 +599,12 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await this.announceCreated({ sessionId: targetId, handle: target, source: 'fork' }); return target; } catch (error) { - if (targetId !== undefined && target !== undefined) { - await this.rollbackFailedFork(targetId, target, targetSessionDir); + if (targetId !== undefined && target !== undefined && targetClaim !== undefined) { + await this.rollbackFailedFork(targetId, target, targetSessionDir, targetClaim); } throw error; } finally { + targetClaim?.release(); quiesce?.dispose(); } } @@ -788,6 +819,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/test/app/sessionLifecycle/sessionLifecycle.test.ts b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts index 17c975a7bc..cbe6748de3 100644 --- a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts +++ b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts @@ -783,37 +783,113 @@ 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('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; + }, }), ]); - let settled = false; - const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }).then(() => { - settled = true; + const completionOrder: string[] = []; + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }).then((handle) => { + completionOrder.push('create'); + return handle; + }); + await mcpStarted.promise; + const resume = svc.resume('s1').then((handle) => { + completionOrder.push('resume'); + return handle; }); - await tick(); - expect(settled).toBe(false); + expect(svc.get('s1')).toBeUndefined(); + expect(svc.list()).toEqual([]); - resolveMcpReady?.(); - await create; - expect(settled).toBe(true); + mcpReady.resolve(); + const created = await create; + await expect(resume).resolves.toBe(created); + expect(completionOrder).toEqual(['create', 'resume']); + 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('honors close while session initialization is still in flight', 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; + await svc.close('s1'); + expect(closed).toEqual(['s1']); + + const rejected = expect(create).rejects.toThrow(/disposed/); + firstMcp.resolve(); + await rejected; + const retried = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); + expect(svc.get('s1')).toBe(retried); }); - it('allows a same-id create retry after MCP initialization fails', async () => { + 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: () => { @@ -829,16 +905,21 @@ describe('SessionLifecycleService', () => { const firstCreate = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); await firstMcpStarted.promise; - const failedHandle = svc.get('s1'); + 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(); - expect(() => failedHandle?.accessor.get(ISessionContext)).toThrow(/disposed/); 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('allows a same-id resume retry after its creation hook fails', async () => { @@ -862,38 +943,6 @@ describe('SessionLifecycleService', () => { hook.dispose(); }); - it('preserves a concurrent replacement when an older materialization fails', async () => { - const firstMcpStarted = deferred(); - const firstMcp = deferred(); - const failure = new Error('older MCP initialization failed'); - let attempts = 0; - const svc = build([ - 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; - const failedHandle = svc.get('s1'); - const replacement = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); - const rejected = expect(firstCreate).rejects.toBe(failure); - firstMcp.reject(failure); - await rejected; - - expect(svc.get('s1')).toBe(replacement); - expect(svc.list()).toEqual([replacement]); - expect(() => failedHandle?.accessor.get(ISessionContext)).toThrow(/disposed/); - }); - it('allows a same-id create retry when failed-session agent cleanup rejects', async () => { const flushStarted = deferred(); const flush = deferred(); @@ -933,7 +982,6 @@ describe('SessionLifecycleService', () => { const creation = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); await flushStarted.promise; - const failedHandle = svc.get('s1'); const rejected = expect(creation).rejects.toBe(initializationFailure); flush.reject(initializationFailure); await rejected; @@ -941,7 +989,6 @@ describe('SessionLifecycleService', () => { expect(removed).toHaveLength(2); expect(removed).toEqual(expect.arrayContaining(['main', 'subagent'])); expect(svc.get('s1')).toBeUndefined(); - expect(() => failedHandle?.accessor.get(ISessionContext)).toThrow(/disposed/); const retried = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); expect(svc.get('s1')).toBe(retried); @@ -1349,71 +1396,96 @@ describe('SessionLifecycleService', () => { expect(replacement.accessor.get(ISessionContext).sessionId).toBe('dst'); }); - it('delays a same-id replacement until failed fork directory rollback finishes', async () => { + 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 cleanupStarted = deferred(); - const releaseCleanup = deferred(); - const replacementReachedRegistration = deferred(); - const cleanupAgent = { - id: 'main', - kind: LifecycleScope.Agent, - accessor: { get: () => ({}) }, - dispose: () => {}, - } as unknown as IAgentScopeHandle; + 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(ISessionMetadata, { - ...metadataStub(), - read: () => - Promise.resolve({ - agents: { main: { homedir: join(srcDir, 'agents', 'main') } }, - } as never), - }), - stubPair(IAgentLifecycleService, { - ...agentLifecycleStub(), - list: () => [cleanupAgent], - remove: () => { - cleanupStarted.resolve(); - return releaseCleanup.promise; - }, - }), - stubPair(ISessionWorkspaceContext, { - _serviceBrand: undefined, - workDir: '/tmp/proj', - additionalDirs: [], - setWorkDir: () => {}, - setAdditionalDirs: () => { - replacementReachedRegistration.resolve(); + 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; }, - resolve: (path: string) => path, - isWithin: () => true, - assertAllowed: (path: string) => path, - addAdditionalDir: () => {}, - removeAdditionalDir: () => {}, }), ]); - await svc.create({ sessionId: 'src', workDir: '/tmp/proj' }); + 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 cleanupStarted.promise; - const forkRejected = expect(forked).rejects.toThrow('not implemented'); - const replacement = svc.create({ - sessionId: 'dst', - workDir: '/tmp/proj', - additionalDirs: ['/tmp/extra'], + 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 replacementReachedRegistration.promise; + 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', + ); - expect(svc.get('dst')).toBeUndefined(); + 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, + }); - releaseCleanup.resolve(); + releaseRollback.resolve(); await forkRejected; - const replacementHandle = await replacement; - expect(svc.get('dst')).toBe(replacementHandle); + 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 () => { From 9e4679d44ffbe33e3a8575ccab2f56a7f503a71d Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 14 Jul 2026 07:32:42 +0800 Subject: [PATCH 4/8] fix(agent-core-v2): separate session publication from teardown --- .../app/sessionLifecycle/sessionLifecycle.ts | 29 +- .../sessionLifecycleService.ts | 293 ++++++++++--- .../sessionLifecycle/sessionLifecycle.test.ts | 402 +++++++++++++++++- 3 files changed, 624 insertions(+), 100 deletions(-) diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts index 68c48eb887..c7a7ed302d 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycle.ts @@ -99,26 +99,29 @@ export interface ISessionLifecycleService { create(opts: CreateSessionOptions): Promise; /** * Return the live handle for `sessionId`, or `undefined` when it is not open. - * A session whose create, fork, or cold {@link resume} initialization is - * still in flight is intentionally NOT returned. Its metadata, agents, or - * replay may be incomplete. Callers that must obtain the handle should - * `await resume(sessionId)` instead. This invisibility is a service - * invariant: every read path (`get` / {@link list} / {@link resume}) agrees - * an initializing 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 - * initializing 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 2b5b25ec4b..470028f9b4 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -14,8 +14,10 @@ * session is also appended to the shared `session_index.jsonl` so v1 clients * (TUI, export) can discover sessions created by the v2 engine. * Failed materialization, creation, resume, and fork attempts release only the - * Session handle they created. Per-session initialization claims prevent - * overlapping create, resume, and fork attempts from sharing persistence. + * Session handle they created. 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'; @@ -97,10 +99,21 @@ type MaterializeSessionOptions = Omit & { }; 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(); @@ -117,12 +130,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( @@ -155,10 +162,12 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec const main = await ensureMainAgent(handle); await main.accessor.get(IAgentPlanService).enter(); } - await this.announceCreated({ sessionId, handle, source: 'startup' }); + this.publishSessionHandle(sessionId, handle, claim); + await this.announceCreated({ sessionId, handle, source: 'startup' }, claim); + this.assertSessionHandleOwned(sessionId, handle, claim); return handle; } catch (error) { - if (handle !== undefined) { + if (handle !== undefined && !claim.closing) { await this.disposeFailedSession(sessionId, handle); } throw error; @@ -222,10 +231,12 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec handle.accessor.get(ISessionExternalHooksService); return handle; } catch (error) { - if (opts.forkDestination === true) { - await this.rollbackFailedFork(opts.sessionId, handle, ctx.sessionDir, opts.claim); - } else { - await this.disposeFailedSession(opts.sessionId, handle); + if (!opts.claim.closing) { + if (opts.forkDestination === true) { + await this.rollbackFailedFork(opts.sessionId, handle, ctx.sessionDir, opts.claim); + } else { + await this.disposeFailedSession(opts.sessionId, handle); + } } throw error; } @@ -242,8 +253,12 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await this.appendLogStore.flush(); } - private async announceCreated(event: SessionCreatedEvent): Promise { + private async announceCreated( + event: SessionCreatedEvent, + claim: SessionInitializationClaim, + ): Promise { await this.hooks.onDidCreateSession.run(event); + 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 @@ -253,23 +268,31 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } get(sessionId: string): ISessionScopeHandle | undefined { - if (this.initializations.has(sessionId) || 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 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; - if (!this.initializations.has(sessionId)) { - const live = this.sessions.get(sessionId); - if (live !== undefined) return Promise.resolve(live); - } const promise = this.doResume(sessionId) .catch((error: unknown) => { this.telemetry.track2('session_load_failed', { @@ -290,7 +313,9 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec while (true) { const initialization = this.initializations.get(sessionId); if (initialization !== undefined) { - await initialization.settled; + if (initialization.closing) return undefined; + if (initialization.published) return this.sessions.get(sessionId); + await initialization.publishedOrSettled; continue; } @@ -306,14 +331,22 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec const racedInitialization = this.initializations.get(sessionId); if (racedInitialization !== undefined) { - await racedInitialization.settled; + 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; - const claim = this.claimInitialization(sessionId); + 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({ @@ -331,10 +364,12 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec const records = mainWireRecord.getRecords() as readonly PersistedRecord[]; await main.accessor.get(IAgentWireService).replay(...records); } - await this.announceCreated({ sessionId, handle, source: 'resume' }); + 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) { + if (handle !== undefined && !claim.closing) { await this.disposeFailedSession(sessionId, handle); } throw error; @@ -345,41 +380,65 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } 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.initializations.has(id) && !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 { @@ -427,20 +486,65 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec if (this.sessions.has(sessionId) || this.initializations.has(sessionId)) { throw sessionAlreadyExistsError(sessionId); } - let released = false; + 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, - release: () => { - if (released) return; - released = true; - if (this.initializations.get(sessionId) === claim) { - this.initializations.delete(sessionId); + markPublished: () => { + if (writerReleased || published) return; + published = true; + resolvePublishedOrSettled(); + }, + startLifecycleAction: () => { + if (claimSettled || writerReleased || lifecycleActionStarted) { + return false; } - settle(); + 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); @@ -458,6 +562,55 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec this.sessions.set(sessionId, handle); } + private publishSessionHandle( + sessionId: string, + handle: ISessionScopeHandle, + claim: SessionInitializationClaim, + ): void { + this.assertSessionHandleOwned(sessionId, handle, claim); + 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 rollbackFailedFork( sessionId: string, handle: ISessionScopeHandle, @@ -591,15 +744,25 @@ 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) { - if (targetId !== undefined && target !== undefined && targetClaim !== undefined) { + if ( + targetId !== undefined && + target !== undefined && + targetClaim !== undefined && + !targetClaim.closing + ) { await this.rollbackFailedFork(targetId, target, targetSessionDir, targetClaim); } throw error; 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 cbe6748de3..45b8a0c8ce 100644 --- a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts +++ b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts @@ -396,10 +396,6 @@ function deferred(): { return { promise, resolve, reject }; } -function tick(): Promise { - return new Promise((resolve) => setTimeout(resolve, 0)); -} - class NoopSessionExternalHooksService implements ISessionExternalHooksService { declare readonly _serviceBrand: undefined; } @@ -783,6 +779,209 @@ describe('SessionLifecycleService', () => { expect(recordedSessionHookEvents).toEqual(['create:startup:s1', 'close:exit:s1']); }); + 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('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(); @@ -796,24 +995,17 @@ describe('SessionLifecycleService', () => { }), ]); - const completionOrder: string[] = []; - const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }).then((handle) => { - completionOrder.push('create'); - return handle; - }); + const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); await mcpStarted.promise; - const resume = svc.resume('s1').then((handle) => { - completionOrder.push('resume'); - return handle; - }); + const resume = svc.resume('s1'); expect(svc.get('s1')).toBeUndefined(); expect(svc.list()).toEqual([]); mcpReady.resolve(); + const resumed = await resume; const created = await create; - await expect(resume).resolves.toBe(created); - expect(completionOrder).toEqual(['create', 'resume']); + expect(resumed).toBe(created); expect(svc.get('s1')).toBe(created); expect(svc.list()).toEqual([created]); }); @@ -843,7 +1035,7 @@ describe('SessionLifecycleService', () => { expect(svc.get('s1')).toBeUndefined(); }); - it('honors close while session initialization is still in flight', async () => { + it('honors close after an in-flight session reaches publication', async () => { const firstMcpStarted = deferred(); const firstMcp = deferred(); let attempts = 0; @@ -865,16 +1057,122 @@ describe('SessionLifecycleService', () => { const create = svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); await firstMcpStarted.promise; - await svc.close('s1'); - expect(closed).toEqual(['s1']); - - const rejected = expect(create).rejects.toThrow(/disposed/); + 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']); const retried = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); expect(svc.get('s1')).toBe(retried); }); + 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(); @@ -995,6 +1293,7 @@ describe('SessionLifecycleService', () => { }); 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; @@ -1003,12 +1302,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. @@ -1367,6 +1669,62 @@ describe('SessionLifecycleService', () => { 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(); From 6d8c880dd885dccf59d3673f804653be0ee2510f Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 14 Jul 2026 07:58:37 +0800 Subject: [PATCH 5/8] fix(agent-core-v2): wait for stable fork sources --- .../sessionLifecycleService.ts | 11 +- .../sessionLifecycle/sessionLifecycle.test.ts | 162 ++++++++++++++++++ 2 files changed, 172 insertions(+), 1 deletion(-) diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts index 470028f9b4..50f702b531 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -611,6 +611,15 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec return { handle, claim: initialization }; } + private async forkSourceHandle(sessionId: string): Promise { + let initialization = this.initializations.get(sessionId); + while (initialization !== undefined) { + await initialization.settled; + initialization = this.initializations.get(sessionId); + } + return this.sessions.get(sessionId); + } + private async rollbackFailedFork( sessionId: string, handle: ISessionScopeHandle, @@ -633,7 +642,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec // 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`); 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 45b8a0c8ce..ac9f2ec300 100644 --- a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts +++ b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts @@ -1528,6 +1528,168 @@ describe('SessionLifecycleService', () => { }); } + 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([ From 20e75aed95522e7c41d060be4f950cd9e74a0849 Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 14 Jul 2026 09:34:05 +0800 Subject: [PATCH 6/8] fix(agent-core-v2): roll back failed fresh sessions --- .../app/sessionIndex/sessionIndexService.ts | 27 +- .../sessionLifecycleService.ts | 145 ++++++- .../sessionMetadata/sessionMetadata.ts | 6 +- .../sessionMetadata/sessionMetadataService.ts | 17 +- .../app/sessionExport/sessionExport.test.ts | 1 + .../app/sessionIndex/sessionIndex.test.ts | 45 ++- .../sessionLifecycle/sessionLifecycle.test.ts | 371 +++++++++++++++++- .../sessionMetadata/sessionMetadata.test.ts | 70 ++++ packages/agent-core-v2/test/tool/tool.test.ts | 1 + 9 files changed, 633 insertions(+), 50 deletions(-) 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/sessionLifecycleService.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts index 50f702b531..a60a6005f3 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -13,8 +13,10 @@ * roots are remembered through `workspaceRegistry`. On create / fork the * session is also appended to the shared `session_index.jsonl` so v1 clients * (TUI, export) can discover sessions created by the v2 engine. - * Failed materialization, creation, resume, and fork attempts release only the - * Session handle they created. Per-session initialization claims keep writers + * 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. @@ -22,7 +24,7 @@ import { randomUUID } from 'node:crypto'; -import { join } from 'pathe'; +import { dirname, join } from 'pathe'; import { ulid } from 'ulid'; import { InstantiationType } from '#/_base/di/extensions'; @@ -96,8 +98,14 @@ type MaterializeSessionOptions = Omit & { 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; @@ -154,9 +162,15 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec async create(opts: CreateSessionOptions): Promise { const sessionId = opts.sessionId ?? createSessionId(); const claim = this.claimInitialization(sessionId); + const directoryOwnership: SessionDirectoryOwnership = { owned: false }; let handle: ISessionScopeHandle | undefined; try { - handle = await this.materializeSession({ ...opts, sessionId, claim }); + 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); @@ -167,8 +181,18 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec this.assertSessionHandleOwned(sessionId, handle, claim); return handle; } catch (error) { - if (handle !== undefined && !claim.closing) { - await this.disposeFailedSession(sessionId, handle); + const failedHandle = handle; + if (failedHandle !== undefined && !claim.closing) { + await throwAfterCleanup(error, () => + directoryOwnership.owned + ? this.rollbackOwnedSessionDirectory( + sessionId, + failedHandle, + directoryOwnership.sessionDir, + claim, + ) + : this.disposeFailedSession(sessionId, failedHandle), + ); } throw error; } finally { @@ -211,15 +235,25 @@ 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; + 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); @@ -231,11 +265,25 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec handle.accessor.get(ISessionExternalHooksService); return handle; } catch (error) { + const failedHandle = handle; if (!opts.claim.closing) { - if (opts.forkDestination === true) { - await this.rollbackFailedFork(opts.sessionId, handle, ctx.sessionDir, opts.claim); - } else { - await this.disposeFailedSession(opts.sessionId, handle); + 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; @@ -477,6 +525,9 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } catch {} } } catch {} + try { + await handle.accessor.get(ISessionMetadata).whenIdle(); + } catch {} try { handle.dispose(); } catch {} @@ -620,7 +671,17 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec return this.sessions.get(sessionId); } - private async rollbackFailedFork( + 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, @@ -634,7 +695,17 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec this.sessions.delete(sessionId); } await this.disposeFailedSession(sessionId, handle); - await this.hostFs.remove(sessionDir).catch(() => {}); + 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 { @@ -772,7 +843,17 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec targetClaim !== undefined && !targetClaim.closing ) { - await this.rollbackFailedFork(targetId, target, targetSessionDir, targetClaim); + const failedTargetId = targetId; + const failedTarget = target; + const failedTargetClaim = targetClaim; + await throwAfterCleanup(error, () => + this.rollbackOwnedSessionDirectory( + failedTargetId, + failedTarget, + targetSessionDir, + failedTargetClaim, + ), + ); } throw error; } finally { @@ -979,6 +1060,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`). 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 ac9f2ec300..51b39b5118 100644 --- a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts +++ b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts @@ -20,6 +20,7 @@ import { Disposable } from '#/_base/di/lifecycle'; import { ILogService } from '#/_base/log/log'; import { type IAgentScopeHandle, + type ISessionScopeHandle, LifecycleScope, type ScopeSeed, _clearScopedRegistryForTests, @@ -76,15 +77,8 @@ import { recordingTelemetry, type TelemetryRecord } from '../../app/telemetry/st 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 { @@ -126,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(), @@ -282,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, @@ -429,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, @@ -476,7 +493,7 @@ describe('SessionLifecycleService', () => { function build(extra: ScopeSeed = []): ISessionLifecycleService { host = createScopedTestHost([ - stubPair(IBootstrapService, bootstrapStub()), + stubPair(IBootstrapService, bootstrapStub(defaultRoot)), stubPair(ISessionMetadata, metadataStub()), stubPair(IHostEnvironment, hostEnvironmentStub()), stubPair(ISessionSkillCatalog, skillCatalogStub()), @@ -484,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()), @@ -542,7 +561,7 @@ describe('SessionLifecycleService', () => { key: 'session_index.jsonl', record: { sessionId: 's1', - sessionDir: `/tmp/sessions/${workspaceId}/s1`, + sessionDir: join(defaultRoot, 'sessions', workspaceId, 's1'), workDir: '/tmp/proj', }, }, @@ -631,7 +650,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 () => { @@ -1035,6 +1054,327 @@ describe('SessionLifecycleService', () => { 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 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(); + }); + + 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(); + + 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(); @@ -1065,8 +1405,9 @@ describe('SessionLifecycleService', () => { await rejected; await close; expect(closed).toEqual(['s1']); - const retried = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' }); - expect(svc.get('s1')).toBe(retried); + 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 () => { 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, From 75502ee8da80b2d6fbc3bd63f333e2ec5fa6df10 Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 14 Jul 2026 09:59:12 +0800 Subject: [PATCH 7/8] fix(agent-core-v2): preserve published session handles --- .../sessionLifecycleService.ts | 15 +++-- .../sessionLifecycle/sessionLifecycle.test.ts | 67 +++++++++++++++++-- 2 files changed, 74 insertions(+), 8 deletions(-) diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts index a60a6005f3..b1b876eb96 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -182,7 +182,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec return handle; } catch (error) { const failedHandle = handle; - if (failedHandle !== undefined && !claim.closing) { + if (failedHandle !== undefined && !claim.published && !claim.closing) { await throwAfterCleanup(error, () => directoryOwnership.owned ? this.rollbackOwnedSessionDirectory( @@ -305,14 +305,19 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec event: SessionCreatedEvent, claim: SessionInitializationClaim, ): Promise { - await this.hooks.onDidCreateSession.run(event); + 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 { @@ -417,7 +422,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec this.assertSessionHandleOwned(sessionId, handle, claim); return handle; } catch (error) { - if (handle !== undefined && !claim.closing) { + if (handle !== undefined && !claim.published && !claim.closing) { await this.disposeFailedSession(sessionId, handle); } throw error; @@ -619,6 +624,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec claim: SessionInitializationClaim, ): void { this.assertSessionHandleOwned(sessionId, handle, claim); + handle.accessor.get(ISessionActivityKernel).markActive(); claim.markPublished(); this.assertSessionHandleOwned(sessionId, handle, claim); } @@ -841,6 +847,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec targetId !== undefined && target !== undefined && targetClaim !== undefined && + !targetClaim.published && !targetClaim.closing ) { const failedTargetId = targetId; 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 51b39b5118..869bf30007 100644 --- a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts +++ b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts @@ -818,6 +818,38 @@ describe('SessionLifecycleService', () => { 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[] = []; @@ -1561,7 +1593,7 @@ describe('SessionLifecycleService', () => { expect(workspaceTouches).toBe(2); }); - it('allows a same-id resume retry after its creation hook fails', async () => { + 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([ @@ -1575,10 +1607,11 @@ describe('SessionLifecycleService', () => { }); await expect(svc.resume('s1')).rejects.toBe(failure); - expect(svc.get('s1')).toBeUndefined(); + const published = svc.get('s1'); - const retried = await svc.resume('s1'); - expect(svc.get('s1')).toBe(retried); + expect(published).toBeDefined(); + await expect(svc.resume('s1')).resolves.toBe(published); + expect(attempts).toBe(1); hook.dispose(); }); @@ -1869,6 +1902,32 @@ describe('SessionLifecycleService', () => { }); } + 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(); From 1783b485f90ea7373234a78b2ae7497b9f18abb8 Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 14 Jul 2026 10:19:36 +0800 Subject: [PATCH 8/8] fix(agent-core-v2): allow hooks to fork published sessions --- .../sessionLifecycleService.ts | 7 ++++++- .../sessionLifecycle/sessionLifecycle.test.ts | 20 +++++++++++++++++++ 2 files changed, 26 insertions(+), 1 deletion(-) diff --git a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts index b1b876eb96..046a18f5cf 100644 --- a/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/app/sessionLifecycle/sessionLifecycleService.ts @@ -671,7 +671,12 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec private async forkSourceHandle(sessionId: string): Promise { let initialization = this.initializations.get(sessionId); while (initialization !== undefined) { - await initialization.settled; + 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); 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 869bf30007..d152b3a78c 100644 --- a/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts +++ b/packages/agent-core-v2/test/app/sessionLifecycle/sessionLifecycle.test.ts @@ -1902,6 +1902,26 @@ 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;