From abcbe5ea136bd264ebacdee6b16df2b4be56dedd Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 14:27:16 +0800 Subject: [PATCH 01/10] fix(agent-core-v2): cache workspace alias resolution across calls resolveAliasIds re-read the workspace catalog and the whole session index from disk on every call, so the by_workspace grouping loop and the per-workspace session counts paid repeated full-file reads per workspace per request (~2.4s per 50-group page at 1.1k workspaces, 23 pages serially during a client startup drain). Cache both files as precomputed snapshots (by-id map plus a root-key -> alias ids index) invalidated by the storage watch events, which cover atomic rewrites and cross-process writes; storage backends without watch fall back to reading through. The first resolution primes the workspace merge via IWorkspaceService.list() so the cached catalog matches what WorkspaceService.get() would have returned. --- .changeset/workspace-aliases-cache.md | 5 + .../app/workspace/fileWorkspacePersistence.ts | 4 +- .../workspaceAliasesService.ts | 118 ++++++++++++++++-- .../workspaceAliasesService.test.ts | 78 +++++++++++- 4 files changed, 189 insertions(+), 16 deletions(-) create mode 100644 .changeset/workspace-aliases-cache.md diff --git a/.changeset/workspace-aliases-cache.md b/.changeset/workspace-aliases-cache.md new file mode 100644 index 00000000000..44cbea4cde8 --- /dev/null +++ b/.changeset/workspace-aliases-cache.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Fix slow session and workspace listing on installs with many workspaces. diff --git a/packages/agent-core-v2/src/app/workspace/fileWorkspacePersistence.ts b/packages/agent-core-v2/src/app/workspace/fileWorkspacePersistence.ts index 920412fe897..14873711495 100644 --- a/packages/agent-core-v2/src/app/workspace/fileWorkspacePersistence.ts +++ b/packages/agent-core-v2/src/app/workspace/fileWorkspacePersistence.ts @@ -12,8 +12,8 @@ import { } from './workspacePersistence'; const WORKSPACE_CATALOG_VERSION = 1; -const WORKSPACE_CATALOG_SCOPE = ''; -const WORKSPACE_CATALOG_KEY = 'workspaces.json'; +export const WORKSPACE_CATALOG_SCOPE = ''; +export const WORKSPACE_CATALOG_KEY = 'workspaces.json'; export class FileWorkspacePersistence implements IWorkspacePersistence { declare readonly _serviceBrand: undefined; diff --git a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts index 8fc84bcccd9..a1d6d515ba2 100644 --- a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts +++ b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts @@ -1,34 +1,130 @@ import { LifecycleScope } from '#/app/scopes'; +import { Disposable } from '#/_base/di/lifecycle'; import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; -import { IWorkspaceService } from '#/app/workspace/workspace'; +import { encodeWorkDirKey, workspaceRootKey } from '#/_base/utils/workdir-slug'; +import { IWorkspaceService, type Workspace } from '#/app/workspace/workspace'; import { - collectAliasIds, readSessionIndexEntries, + SESSION_INDEX_KEY, + SESSION_INDEX_SCOPE, } from '#/app/workspace/workspaceAlias'; +import { + WORKSPACE_CATALOG_KEY, + WORKSPACE_CATALOG_SCOPE, +} from '#/app/workspace/fileWorkspacePersistence'; import { IWorkspacePersistence } from '#/app/workspace/workspacePersistence'; import { IFileSystemStorageService } from '#/persistence/interface/storage'; import { IWorkspaceAliases } from './workspaceAliases'; -export class WorkspaceAliasesService implements IWorkspaceAliases { +interface CatalogSnapshot { + readonly byId: ReadonlyMap; + readonly idsByRootKey: ReadonlyMap; +} + +interface SessionIndexSnapshot { + readonly idsByRootKey: ReadonlyMap; +} + +function rootKeyIndex( + items: readonly T[], + rootOf: (item: T) => string, + idOf: (item: T) => string, +): Map { + const map = new Map(); + for (const item of items) { + const key = workspaceRootKey(rootOf(item)); + const id = idOf(item); + const bucket = map.get(key); + if (bucket === undefined) { + map.set(key, [id]); + } else if (!bucket.includes(id)) { + bucket.push(id); + } + } + return map; +} + +export class WorkspaceAliasesService extends Disposable implements IWorkspaceAliases { declare readonly _serviceBrand: undefined; + private catalogCache: CatalogSnapshot | undefined; + private sessionIndexCache: SessionIndexSnapshot | undefined; + private readonly cacheable: boolean; + private catalogMergePrimed = false; + constructor( @IWorkspaceService private readonly workspaces: IWorkspaceService, @IWorkspacePersistence private readonly store: IWorkspacePersistence, @IFileSystemStorageService private readonly storage: IFileSystemStorageService, - ) {} + ) { + super(); + const watchCatalog = this.storage.watch?.(WORKSPACE_CATALOG_SCOPE, WORKSPACE_CATALOG_KEY); + const watchSessionIndex = this.storage.watch?.(SESSION_INDEX_SCOPE, SESSION_INDEX_KEY); + this.cacheable = watchCatalog !== undefined && watchSessionIndex !== undefined; + if (watchCatalog !== undefined) { + this._register( + watchCatalog(() => { + this.catalogCache = undefined; + }), + ); + } + if (watchSessionIndex !== undefined) { + this._register( + watchSessionIndex(() => { + this.sessionIndexCache = undefined; + }), + ); + } + } async resolveAliasIds(id: string): Promise { - const entry = await this.workspaces.get(id); + const catalog = await this.catalog(); + const entry = catalog.byId.get(id); if (entry === undefined) return [id]; - const catalog = (await this.store.load()) ?? { workspaces: [], deletedIds: [] }; - return collectAliasIds( - catalog.workspaces, - await readSessionIndexEntries(this.storage), - entry.root, - ); + const rootKey = workspaceRootKey(entry.root); + const index = await this.sessionIndex(); + const fromCatalog = catalog.idsByRootKey.get(rootKey); + const fromIndex = index.idsByRootKey.get(rootKey); + if (fromCatalog === undefined) return fromIndex ?? [id]; + if (fromIndex === undefined) return fromCatalog; + const merged = [...fromCatalog]; + for (const alias of fromIndex) { + if (!merged.includes(alias)) merged.push(alias); + } + return merged; + } + + private async catalog(): Promise { + if (this.catalogCache !== undefined) return this.catalogCache; + if (!this.catalogMergePrimed) { + await this.workspaces.list(); + this.catalogMergePrimed = true; + } + const workspaces = (await this.store.load())?.workspaces ?? []; + const snapshot: CatalogSnapshot = { + byId: new Map(workspaces.map((ws) => [ws.id, ws] as const)), + idsByRootKey: rootKeyIndex( + workspaces, + (ws) => ws.root, + (ws) => ws.id, + ), + }; + if (this.cacheable) this.catalogCache = snapshot; + return snapshot; + } + + private async sessionIndex(): Promise { + if (this.sessionIndexCache !== undefined) return this.sessionIndexCache; + const entries = await readSessionIndexEntries(this.storage); + const snapshot: SessionIndexSnapshot = { + idsByRootKey: rootKeyIndex(entries, (entry) => entry.workDir, (entry) => + encodeWorkDirKey(entry.workDir), + ), + }; + if (this.cacheable) this.sessionIndexCache = snapshot; + return snapshot; } } diff --git a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts index edec7f532c1..149e58f085f 100644 --- a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts +++ b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts @@ -1,4 +1,4 @@ -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { promises as fsp } from 'node:fs'; import os from 'node:os'; @@ -70,8 +70,10 @@ describe('WorkspaceAliasesService (file-backed)', () => { await fsp.rm(homeDir, { recursive: true, force: true }); }); - function build(hostFs: IHostFileSystem = new HostFileSystem()): IWorkspaceAliases { - const fileStorage = new FileStorageService(homeDir); + function build( + hostFs: IHostFileSystem = new HostFileSystem(), + fileStorage: FileStorageService = new FileStorageService(homeDir), + ): IWorkspaceAliases { const host = createScopedTestHost([ stubPair(IFileSystemStorageService, fileStorage), stubPair(IAtomicDocumentStore, new JsonAtomicDocumentStore(fileStorage)), @@ -168,4 +170,74 @@ describe('WorkspaceAliasesService (file-backed)', () => { ]); expect(await aliases.resolveAliasIds(id)).toEqual([id]); }); + + it('resolveAliasIds reuses the loaded catalog and session index across calls', async () => { + class CountingStorage extends FileStorageService { + reads = 0; + override async read(scope: string, key: string): Promise { + this.reads += 1; + return super.read(scope, key); + } + } + const root = join(homeDir, 'proj'); + const id = encodeWorkDirKey(root); + await writeWorkspacesJson({ + [id]: { + root, + name: 'proj', + created_at: '2026-01-01T00:00:00.000Z', + last_opened_at: '2026-01-01T00:00:00.000Z', + }, + }); + await seedSessionIndex([{ sessionId: 's1', sessionDir: 'sessions/a/s1', workDir: root }]); + const storage = new CountingStorage(homeDir); + const aliases = build(undefined, storage); + + await aliases.resolveAliasIds(id); + const readsAfterWarm = storage.reads; + await aliases.resolveAliasIds(id); + await aliases.resolveAliasIds('wd_missing_000000000000'); + expect(storage.reads).toBe(readsAfterWarm); + }); + + it('resolveAliasIds picks up catalog and session index changes', async () => { + const entry = (root: string): PersistedWorkspaceEntry => ({ + root, + name: 'proj', + created_at: '2026-01-01T00:00:00.000Z', + last_opened_at: '2026-01-01T00:00:00.000Z', + }); + const typedRoot = 'C:\\Users\\Foo\\Proj'; + const typedId = encodeWorkDirKey(typedRoot); + const legacyId = 'wd_proj_deadbeef0002'; + await writeWorkspacesJson({ [typedId]: entry(typedRoot) }); + const aliases = build(); + expect(await aliases.resolveAliasIds(typedId)).toEqual([typedId]); + + await writeWorkspacesJson({ + [typedId]: entry(typedRoot), + [legacyId]: entry('c:\\users\\foo\\proj'), + }); + await vi.waitFor( + async () => { + expect((await aliases.resolveAliasIds(typedId)).toSorted()).toEqual( + [legacyId, typedId].toSorted(), + ); + }, + { timeout: 5000 }, + ); + + const indexOnlyId = encodeWorkDirKey('c:\\Users\\Foo\\Proj'); + await seedSessionIndex([ + { sessionId: 's1', sessionDir: 'sessions/a/s1', workDir: 'c:\\Users\\Foo\\Proj' }, + ]); + await vi.waitFor( + async () => { + expect((await aliases.resolveAliasIds(typedId)).toSorted()).toEqual( + [indexOnlyId, legacyId, typedId].toSorted(), + ); + }, + { timeout: 5000 }, + ); + }); }); From d9c829735706535246c48d2be38e83d6ff630377 Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 14:31:50 +0800 Subject: [PATCH 02/10] chore: drop the changeset; the user-facing entry ships with the app changelog --- .changeset/workspace-aliases-cache.md | 5 ----- 1 file changed, 5 deletions(-) delete mode 100644 .changeset/workspace-aliases-cache.md diff --git a/.changeset/workspace-aliases-cache.md b/.changeset/workspace-aliases-cache.md deleted file mode 100644 index 44cbea4cde8..00000000000 --- a/.changeset/workspace-aliases-cache.md +++ /dev/null @@ -1,5 +0,0 @@ ---- -"@moonshot-ai/kimi-code": patch ---- - -Fix slow session and workspace listing on installs with many workspaces. From b26483749fa85931bc3b157c5e3747053eb1df50 Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 14:42:00 +0800 Subject: [PATCH 03/10] fix(agent-core-v2): coalesce cold alias snapshot loads and guard publication Concurrent cold resolveAliasIds callers (the /workspaces route fans out per-workspace counts with Promise.all) all passed the cache check before any caller finished loading, re-running the full catalog and session index reads the cache exists to avoid; memoize the in-flight load promise so a cold batch shares one read. Also capture the invalidation generation before each read and publish the snapshot only when it is unchanged, so a mid-read file replacement cannot leave a stale snapshot installed over the watch invalidation. --- .../workspaceAliasesService.ts | 81 +++++++++++++------ .../workspaceAliasesService.test.ts | 39 +++++++-- 2 files changed, 87 insertions(+), 33 deletions(-) diff --git a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts index a1d6d515ba2..bcee9fe2319 100644 --- a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts +++ b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts @@ -51,6 +51,9 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli private catalogCache: CatalogSnapshot | undefined; private sessionIndexCache: SessionIndexSnapshot | undefined; + private catalogPromise: Promise | undefined; + private sessionIndexPromise: Promise | undefined; + private invalidationGeneration = 0; private readonly cacheable: boolean; private catalogMergePrimed = false; @@ -66,6 +69,7 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli if (watchCatalog !== undefined) { this._register( watchCatalog(() => { + this.invalidationGeneration += 1; this.catalogCache = undefined; }), ); @@ -73,6 +77,7 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli if (watchSessionIndex !== undefined) { this._register( watchSessionIndex(() => { + this.invalidationGeneration += 1; this.sessionIndexCache = undefined; }), ); @@ -96,35 +101,59 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli return merged; } - private async catalog(): Promise { - if (this.catalogCache !== undefined) return this.catalogCache; - if (!this.catalogMergePrimed) { - await this.workspaces.list(); - this.catalogMergePrimed = true; + private catalog(): Promise { + if (this.catalogCache !== undefined) return Promise.resolve(this.catalogCache); + this.catalogPromise ??= this.loadCatalog(); + return this.catalogPromise; + } + + private async loadCatalog(): Promise { + try { + if (!this.catalogMergePrimed) { + await this.workspaces.list(); + this.catalogMergePrimed = true; + } + const generation = this.invalidationGeneration; + const workspaces = (await this.store.load())?.workspaces ?? []; + const snapshot: CatalogSnapshot = { + byId: new Map(workspaces.map((ws) => [ws.id, ws] as const)), + idsByRootKey: rootKeyIndex( + workspaces, + (ws) => ws.root, + (ws) => ws.id, + ), + }; + if (this.cacheable && generation === this.invalidationGeneration) { + this.catalogCache = snapshot; + } + return snapshot; + } finally { + this.catalogPromise = undefined; } - const workspaces = (await this.store.load())?.workspaces ?? []; - const snapshot: CatalogSnapshot = { - byId: new Map(workspaces.map((ws) => [ws.id, ws] as const)), - idsByRootKey: rootKeyIndex( - workspaces, - (ws) => ws.root, - (ws) => ws.id, - ), - }; - if (this.cacheable) this.catalogCache = snapshot; - return snapshot; } - private async sessionIndex(): Promise { - if (this.sessionIndexCache !== undefined) return this.sessionIndexCache; - const entries = await readSessionIndexEntries(this.storage); - const snapshot: SessionIndexSnapshot = { - idsByRootKey: rootKeyIndex(entries, (entry) => entry.workDir, (entry) => - encodeWorkDirKey(entry.workDir), - ), - }; - if (this.cacheable) this.sessionIndexCache = snapshot; - return snapshot; + private sessionIndex(): Promise { + if (this.sessionIndexCache !== undefined) return Promise.resolve(this.sessionIndexCache); + this.sessionIndexPromise ??= this.loadSessionIndex(); + return this.sessionIndexPromise; + } + + private async loadSessionIndex(): Promise { + try { + const generation = this.invalidationGeneration; + const entries = await readSessionIndexEntries(this.storage); + const snapshot: SessionIndexSnapshot = { + idsByRootKey: rootKeyIndex(entries, (entry) => entry.workDir, (entry) => + encodeWorkDirKey(entry.workDir), + ), + }; + if (this.cacheable && generation === this.invalidationGeneration) { + this.sessionIndexCache = snapshot; + } + return snapshot; + } finally { + this.sessionIndexPromise = undefined; + } } } diff --git a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts index 149e58f085f..e00a2cc3b42 100644 --- a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts +++ b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts @@ -70,6 +70,14 @@ describe('WorkspaceAliasesService (file-backed)', () => { await fsp.rm(homeDir, { recursive: true, force: true }); }); + class CountingStorage extends FileStorageService { + reads = 0; + override async read(scope: string, key: string): Promise { + this.reads += 1; + return super.read(scope, key); + } + } + function build( hostFs: IHostFileSystem = new HostFileSystem(), fileStorage: FileStorageService = new FileStorageService(homeDir), @@ -172,13 +180,6 @@ describe('WorkspaceAliasesService (file-backed)', () => { }); it('resolveAliasIds reuses the loaded catalog and session index across calls', async () => { - class CountingStorage extends FileStorageService { - reads = 0; - override async read(scope: string, key: string): Promise { - this.reads += 1; - return super.read(scope, key); - } - } const root = join(homeDir, 'proj'); const id = encodeWorkDirKey(root); await writeWorkspacesJson({ @@ -200,6 +201,30 @@ describe('WorkspaceAliasesService (file-backed)', () => { expect(storage.reads).toBe(readsAfterWarm); }); + it('resolveAliasIds coalesces concurrent cold loads', async () => { + const root = join(homeDir, 'proj'); + const id = encodeWorkDirKey(root); + await writeWorkspacesJson({ + [id]: { + root, + name: 'proj', + created_at: '2026-01-01T00:00:00.000Z', + last_opened_at: '2026-01-01T00:00:00.000Z', + }, + }); + await seedSessionIndex([{ sessionId: 's1', sessionDir: 'sessions/a/s1', workDir: root }]); + const storage = new CountingStorage(homeDir); + const aliases = build(undefined, storage); + + await Promise.all([ + aliases.resolveAliasIds(id), + aliases.resolveAliasIds(id), + aliases.resolveAliasIds('wd_missing_000000000000'), + aliases.resolveAliasIds(encodeWorkDirKey(join(homeDir, 'nowhere'))), + ]); + expect(storage.reads).toBeLessThanOrEqual(6); + }); + it('resolveAliasIds picks up catalog and session index changes', async () => { const entry = (root: string): PersistedWorkspaceEntry => ({ root, From e3f7bc3cfddd36ab424b44f43a6c619d815cfb2b Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 14:56:18 +0800 Subject: [PATCH 04/10] fix(agent-core-v2): publish catalog invalidation through the persistence owner A debounced fs watch was the only invalidation channel for the alias catalog snapshot, so an in-process catalog write stayed invisible to resolveAliasIds for up to the watch debounce window while the previous read-through code observed every completed write immediately. IWorkspacePersistence now exposes onDidChange: FileWorkspacePersistence fires it synchronously on save and re-fires the underlying document watch (covering atomic rewrites and cross-process writers), and the aliases service subscribes to it instead of watching raw storage keys. The session index snapshot keeps the filesystem watch, matching the read-side ownership of that file. --- .../app/workspace/fileWorkspacePersistence.ts | 21 ++++++++++++--- .../src/app/workspace/workspacePersistence.ts | 3 +++ .../workspaceAliasesService.ts | 27 +++++++------------ .../workspaceAliasesService.test.ts | 27 +++++++++++++++++++ 4 files changed, 57 insertions(+), 21 deletions(-) diff --git a/packages/agent-core-v2/src/app/workspace/fileWorkspacePersistence.ts b/packages/agent-core-v2/src/app/workspace/fileWorkspacePersistence.ts index 14873711495..dd0e13bc672 100644 --- a/packages/agent-core-v2/src/app/workspace/fileWorkspacePersistence.ts +++ b/packages/agent-core-v2/src/app/workspace/fileWorkspacePersistence.ts @@ -1,6 +1,8 @@ import { LifecycleScope } from '#/app/scopes'; +import { Disposable } from '#/_base/di/lifecycle'; import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; +import { Emitter, type Event } from '#/_base/event'; import { IAtomicDocumentStore } from '#/persistence/interface/atomicDocumentStore'; import type { Workspace } from './workspace'; @@ -12,13 +14,23 @@ import { } from './workspacePersistence'; const WORKSPACE_CATALOG_VERSION = 1; -export const WORKSPACE_CATALOG_SCOPE = ''; -export const WORKSPACE_CATALOG_KEY = 'workspaces.json'; +const WORKSPACE_CATALOG_SCOPE = ''; +const WORKSPACE_CATALOG_KEY = 'workspaces.json'; -export class FileWorkspacePersistence implements IWorkspacePersistence { +export class FileWorkspacePersistence extends Disposable implements IWorkspacePersistence { declare readonly _serviceBrand: undefined; - constructor(@IAtomicDocumentStore private readonly docs: IAtomicDocumentStore) {} + private readonly changeEmitter = this._register(new Emitter()); + readonly onDidChange: Event = this.changeEmitter.event; + + constructor(@IAtomicDocumentStore private readonly docs: IAtomicDocumentStore) { + super(); + this._register( + this.docs.watch(WORKSPACE_CATALOG_SCOPE, WORKSPACE_CATALOG_KEY)(() => { + this.changeEmitter.fire(); + }), + ); + } async load(): Promise { const file = await this.docs.get( @@ -70,6 +82,7 @@ export class FileWorkspacePersistence implements IWorkspacePersistence { deleted_workspace_ids: [...catalog.deletedIds], }; await this.docs.set(WORKSPACE_CATALOG_SCOPE, WORKSPACE_CATALOG_KEY, file); + this.changeEmitter.fire(); } } diff --git a/packages/agent-core-v2/src/app/workspace/workspacePersistence.ts b/packages/agent-core-v2/src/app/workspace/workspacePersistence.ts index 4c96e4d6a3e..f3b941e6231 100644 --- a/packages/agent-core-v2/src/app/workspace/workspacePersistence.ts +++ b/packages/agent-core-v2/src/app/workspace/workspacePersistence.ts @@ -1,4 +1,5 @@ import { createDecorator, type ServiceIdentifier } from '#/_base/di/instantiation'; +import type { Event } from '#/_base/event'; import type { Workspace } from './workspace'; @@ -23,6 +24,8 @@ export interface WorkspaceCatalog { export interface IWorkspacePersistence { readonly _serviceBrand: undefined; + readonly onDidChange: Event; + load(): Promise; save(catalog: WorkspaceCatalog): Promise; } diff --git a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts index bcee9fe2319..f5553580194 100644 --- a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts +++ b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts @@ -9,10 +9,6 @@ import { SESSION_INDEX_KEY, SESSION_INDEX_SCOPE, } from '#/app/workspace/workspaceAlias'; -import { - WORKSPACE_CATALOG_KEY, - WORKSPACE_CATALOG_SCOPE, -} from '#/app/workspace/fileWorkspacePersistence'; import { IWorkspacePersistence } from '#/app/workspace/workspacePersistence'; import { IFileSystemStorageService } from '#/persistence/interface/storage'; @@ -54,7 +50,7 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli private catalogPromise: Promise | undefined; private sessionIndexPromise: Promise | undefined; private invalidationGeneration = 0; - private readonly cacheable: boolean; + private readonly sessionIndexCacheable: boolean; private catalogMergePrimed = false; constructor( @@ -63,17 +59,14 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli @IFileSystemStorageService private readonly storage: IFileSystemStorageService, ) { super(); - const watchCatalog = this.storage.watch?.(WORKSPACE_CATALOG_SCOPE, WORKSPACE_CATALOG_KEY); + this._register( + this.store.onDidChange(() => { + this.invalidationGeneration += 1; + this.catalogCache = undefined; + }), + ); const watchSessionIndex = this.storage.watch?.(SESSION_INDEX_SCOPE, SESSION_INDEX_KEY); - this.cacheable = watchCatalog !== undefined && watchSessionIndex !== undefined; - if (watchCatalog !== undefined) { - this._register( - watchCatalog(() => { - this.invalidationGeneration += 1; - this.catalogCache = undefined; - }), - ); - } + this.sessionIndexCacheable = watchSessionIndex !== undefined; if (watchSessionIndex !== undefined) { this._register( watchSessionIndex(() => { @@ -123,7 +116,7 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli (ws) => ws.id, ), }; - if (this.cacheable && generation === this.invalidationGeneration) { + if (generation === this.invalidationGeneration) { this.catalogCache = snapshot; } return snapshot; @@ -147,7 +140,7 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli encodeWorkDirKey(entry.workDir), ), }; - if (this.cacheable && generation === this.invalidationGeneration) { + if (this.sessionIndexCacheable && generation === this.invalidationGeneration) { this.sessionIndexCache = snapshot; } return snapshot; diff --git a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts index e00a2cc3b42..bf8cc0c1dac 100644 --- a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts +++ b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts @@ -225,6 +225,33 @@ describe('WorkspaceAliasesService (file-backed)', () => { expect(storage.reads).toBeLessThanOrEqual(6); }); + it('resolveAliasIds follows in-process catalog writes synchronously', async () => { + const entry = (root: string): PersistedWorkspaceEntry => ({ + root, + name: 'proj', + created_at: '2026-01-01T00:00:00.000Z', + last_opened_at: '2026-01-01T00:00:00.000Z', + }); + const typedRoot = 'C:\\Users\\Foo\\Proj'; + const typedId = encodeWorkDirKey(typedRoot); + const legacyId = 'wd_proj_deadbeef0002'; + await writeWorkspacesJson({ [typedId]: entry(typedRoot) }); + const aliases = build(); + expect(await aliases.resolveAliasIds(typedId)).toEqual([typedId]); + + const persistence = currentHost!.app.accessor.get(IWorkspacePersistence); + await persistence.save({ + workspaces: [ + { id: typedId, root: typedRoot, name: 'proj', createdAt: 0, lastOpenedAt: 0 }, + { id: legacyId, root: 'c:\\users\\foo\\proj', name: 'proj', createdAt: 0, lastOpenedAt: 0 }, + ], + deletedIds: [], + }); + expect((await aliases.resolveAliasIds(typedId)).toSorted()).toEqual( + [legacyId, typedId].toSorted(), + ); + }); + it('resolveAliasIds picks up catalog and session index changes', async () => { const entry = (root: string): PersistedWorkspaceEntry => ({ root, From 843aad2b008da75a4487965f10a6a4f2e7530d4c Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 15:22:05 +0800 Subject: [PATCH 05/10] fix(agent-core-v2): invalidate session alias snapshots on append-log writes A flushed session_index.jsonl append was invisible to resolveAliasIds for up to the fs-watch debounce window, so a sessions request issued right after a session create could resolve the workspace's aliases from the pre-append snapshot. IAppendLogStore now publishes onDidWrite after each durable flush (append batches and rewrites), and the aliases service drops its session-index snapshot through that event; the raw filesystem watch stays as the channel for cross-process writers. --- .../workspaceAliasesService.ts | 14 +++++++-- .../backends/node-fs/appendLogStore.ts | 13 +++++++-- .../persistence/interface/appendLogStore.ts | 8 +++++ .../workspaceAliasesService.test.ts | 29 +++++++++++++++++++ packages/agent-core-v2/test/harness/agent.ts | 1 + .../agentLifecycle/agentLifecycle.test.ts | 1 + .../test/session/subagent/forkParity.test.ts | 1 + packages/agent-core-v2/test/wire/stubs.ts | 3 ++ 8 files changed, 64 insertions(+), 6 deletions(-) diff --git a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts index f5553580194..6e08d47bec8 100644 --- a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts +++ b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts @@ -10,6 +10,7 @@ import { SESSION_INDEX_SCOPE, } from '#/app/workspace/workspaceAlias'; import { IWorkspacePersistence } from '#/app/workspace/workspacePersistence'; +import { IAppendLogStore } from '#/persistence/interface/appendLogStore'; import { IFileSystemStorageService } from '#/persistence/interface/storage'; import { IWorkspaceAliases } from './workspaceAliases'; @@ -50,13 +51,13 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli private catalogPromise: Promise | undefined; private sessionIndexPromise: Promise | undefined; private invalidationGeneration = 0; - private readonly sessionIndexCacheable: boolean; private catalogMergePrimed = false; constructor( @IWorkspaceService private readonly workspaces: IWorkspaceService, @IWorkspacePersistence private readonly store: IWorkspacePersistence, @IFileSystemStorageService private readonly storage: IFileSystemStorageService, + @IAppendLogStore private readonly appendLogs: IAppendLogStore, ) { super(); this._register( @@ -65,8 +66,15 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli this.catalogCache = undefined; }), ); + this._register( + this.appendLogs.onDidWrite((write) => { + if (write.scope === SESSION_INDEX_SCOPE && write.key === SESSION_INDEX_KEY) { + this.invalidationGeneration += 1; + this.sessionIndexCache = undefined; + } + }), + ); const watchSessionIndex = this.storage.watch?.(SESSION_INDEX_SCOPE, SESSION_INDEX_KEY); - this.sessionIndexCacheable = watchSessionIndex !== undefined; if (watchSessionIndex !== undefined) { this._register( watchSessionIndex(() => { @@ -140,7 +148,7 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli encodeWorkDirKey(entry.workDir), ), }; - if (this.sessionIndexCacheable && generation === this.invalidationGeneration) { + if (generation === this.invalidationGeneration) { this.sessionIndexCache = snapshot; } return snapshot; diff --git a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts index b4eee62e9af..73c4a586684 100644 --- a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts +++ b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts @@ -1,6 +1,7 @@ -import { toDisposable, type IDisposable } from '#/_base/di/lifecycle'; +import { Disposable, toDisposable, type IDisposable } from '#/_base/di/lifecycle'; import { LifecycleScope } from '#/app/scopes'; import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; +import { Emitter, type Event } from '#/_base/event'; import { IFileSystemStorageService } from '#/persistence/interface/storage'; import { @@ -8,6 +9,7 @@ import { IAppendLogStore, type AppendLogOptions, type AppendLogReadOptions, + type AppendLogWrite, } from '#/persistence/interface/appendLogStore'; const textEncoder = new TextEncoder(); @@ -33,12 +35,16 @@ interface LogState { onError?: (error: unknown) => void; } -export class AppendLogStore implements IAppendLogStore { +export class AppendLogStore extends Disposable implements IAppendLogStore { declare readonly _serviceBrand: undefined; private readonly logs = new Map(); + private readonly writeEmitter = this._register(new Emitter()); + readonly onDidWrite: Event = this.writeEmitter.event; - constructor(@IFileSystemStorageService private readonly storage: IFileSystemStorageService) {} + constructor(@IFileSystemStorageService private readonly storage: IFileSystemStorageService) { + super(); + } append(scope: string, key: string, record: R, options?: AppendLogOptions): void { const state = this.state(scope, key); @@ -252,6 +258,7 @@ export class AppendLogStore implements IAppendLogStore { } } if (failure !== undefined) throw failure.error; + this.writeEmitter.fire({ scope, key }); } private async drain(scope: string, key: string, state: LogState): Promise { diff --git a/packages/agent-core-v2/src/persistence/interface/appendLogStore.ts b/packages/agent-core-v2/src/persistence/interface/appendLogStore.ts index a1d78ccb91d..973c0afd6e4 100644 --- a/packages/agent-core-v2/src/persistence/interface/appendLogStore.ts +++ b/packages/agent-core-v2/src/persistence/interface/appendLogStore.ts @@ -1,5 +1,6 @@ import { createDecorator, type ServiceIdentifier } from '#/_base/di/instantiation'; import { type IDisposable } from '#/_base/di/lifecycle'; +import type { Event } from '#/_base/event'; import { StorageError, StorageErrors } from '#/persistence/interface/storage'; @@ -31,9 +32,16 @@ export interface AppendLogReadOptions { readonly onTruncate?: (truncation: AppendLogTruncation) => void; } +export interface AppendLogWrite { + readonly scope: string; + readonly key: string; +} + export interface IAppendLogStore { readonly _serviceBrand: undefined; + readonly onDidWrite: Event; + append(scope: string, key: string, record: R, options?: AppendLogOptions): void; read(scope: string, key: string, options?: AppendLogReadOptions): AsyncIterable; rewrite(scope: string, key: string, records: readonly R[]): Promise; diff --git a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts index bf8cc0c1dac..536c9b6a580 100644 --- a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts +++ b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts @@ -13,8 +13,10 @@ import { createScopedTestHost, stubPair } from '#/_base/di/test'; import { encodeWorkDirKey } from '#/_base/utils/workdir-slug'; import { HostFileSystem } from '#/os/backends/node-local/hostFsService'; import { IHostFileSystem } from '#/os/interface/hostFileSystem'; +import { AppendLogStore } from '#/persistence/backends/node-fs/appendLogStore'; 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 { IFileSystemStorageService } from '#/persistence/interface/storage'; import { IEventService } from '#/app/event/event'; @@ -85,6 +87,7 @@ describe('WorkspaceAliasesService (file-backed)', () => { const host = createScopedTestHost([ stubPair(IFileSystemStorageService, fileStorage), stubPair(IAtomicDocumentStore, new JsonAtomicDocumentStore(fileStorage)), + stubPair(IAppendLogStore, new AppendLogStore(fileStorage)), stubPair(IHostFileSystem, hostFs), stubPair(IEventService, { publish: () => {}, @@ -252,6 +255,32 @@ describe('WorkspaceAliasesService (file-backed)', () => { ); }); + it('resolveAliasIds follows in-process session-index writes synchronously', async () => { + const entry = (root: string): PersistedWorkspaceEntry => ({ + root, + name: 'proj', + created_at: '2026-01-01T00:00:00.000Z', + last_opened_at: '2026-01-01T00:00:00.000Z', + }); + const typedRoot = 'C:\\Users\\Foo\\Proj'; + const typedId = encodeWorkDirKey(typedRoot); + await writeWorkspacesJson({ [typedId]: entry(typedRoot) }); + const aliases = build(); + expect(await aliases.resolveAliasIds(typedId)).toEqual([typedId]); + + const appendLogs = currentHost!.app.accessor.get(IAppendLogStore); + appendLogs.append('', 'session_index.jsonl', { + sessionId: 's9', + sessionDir: 'sessions/s/s9', + workDir: 'c:\\Users\\Foo\\Proj', + }); + await appendLogs.flush(); + const indexOnlyId = encodeWorkDirKey('c:\\Users\\Foo\\Proj'); + expect((await aliases.resolveAliasIds(typedId)).toSorted()).toEqual( + [indexOnlyId, typedId].toSorted(), + ); + }); + it('resolveAliasIds picks up catalog and session index changes', async () => { const entry = (root: string): PersistedWorkspaceEntry => ({ root, diff --git a/packages/agent-core-v2/test/harness/agent.ts b/packages/agent-core-v2/test/harness/agent.ts index 27079579646..7f5ce5a65fb 100644 --- a/packages/agent-core-v2/test/harness/agent.ts +++ b/packages/agent-core-v2/test/harness/agent.ts @@ -976,6 +976,7 @@ function reassertServiceOverrides( class PersistenceAppendLogStore implements IAppendLogStore { declare readonly _serviceBrand: undefined; + readonly onDidWrite: IAppendLogStore['onDidWrite'] = Event.None as IAppendLogStore['onDidWrite']; private readonly history: WireRecord[] = []; private readSeeded = false; diff --git a/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts b/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts index dce5a5c8fd2..034fc1b3e00 100644 --- a/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts +++ b/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts @@ -152,6 +152,7 @@ function recordingAppendLog(initial: readonly WireRecord[] = []): { const state: { rewritten?: readonly WireRecord[] } = {}; const store: IAppendLogStore = { _serviceBrand: undefined, + onDidWrite: Event.None as IAppendLogStore['onDidWrite'], append: (_scope: string, _key: string, record: R) => { const persisted = record as unknown as WireRecord; records.push(persisted); diff --git a/packages/agent-core-v2/test/session/subagent/forkParity.test.ts b/packages/agent-core-v2/test/session/subagent/forkParity.test.ts index 28e0fd7015c..2cb036a1fe3 100644 --- a/packages/agent-core-v2/test/session/subagent/forkParity.test.ts +++ b/packages/agent-core-v2/test/session/subagent/forkParity.test.ts @@ -42,6 +42,7 @@ import { stubFlag } from '../../app/flag/stubs'; class ScopedAppendLogStore implements IAppendLogStore { declare readonly _serviceBrand: undefined; private readonly logs = new Map(); + readonly onDidWrite: IAppendLogStore['onDidWrite'] = Event.None as IAppendLogStore['onDidWrite']; recordsFor(scope: string, key: string): WireRecord[] { return structuredClone(this.logs.get(`${scope}/${key}`) ?? []); diff --git a/packages/agent-core-v2/test/wire/stubs.ts b/packages/agent-core-v2/test/wire/stubs.ts index b081c3b05f2..c2c69cb1e3d 100644 --- a/packages/agent-core-v2/test/wire/stubs.ts +++ b/packages/agent-core-v2/test/wire/stubs.ts @@ -1,6 +1,7 @@ import { SyncDescriptor } from '#/_base/di/descriptors'; import { toDisposable } from '#/_base/di/lifecycle'; import type { ServiceRegistration, TestInstantiationService } from '#/_base/di/test'; +import { Event } from '#/_base/event'; import { ILogService } from '#/_base/log/log'; import { IAgentBlobService } from '#/agent/blob/agentBlobService'; import { AgentRuntimeSet } from '#/agent/runtime/agentRuntimeSet'; @@ -36,6 +37,7 @@ interface TestAgentWireDependencies { const noopLog: IAppendLogStore = { _serviceBrand: undefined, + onDidWrite: Event.None as IAppendLogStore['onDidWrite'], append: () => {}, read: async function* () {}, rewrite: async () => {}, @@ -234,6 +236,7 @@ export function recordingWireLog( ): IAppendLogStore { return { _serviceBrand: undefined, + onDidWrite: Event.None as IAppendLogStore['onDidWrite'], append: (_scope, _key, record) => { records.push(record as WireRecord); onAppend?.(record as WireRecord); From 88d4413269b4636705a8b9f2a082da16edfc512e Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 15:32:45 +0800 Subject: [PATCH 06/10] fix(agent-core-v2): fire append-log write events only after actual writes Once a key has a LogState, every global flush() (WireService flushes after ordinary agent persistence) completed it successfully and fired onDidWrite unconditionally, so idle agent activity kept dropping the alias session-index snapshot and forced full re-reads of an unchanged index. drain() now reports whether it appended anything and the write event fires only when a flush actually persisted a batch or a rewrite. --- .../backends/node-fs/appendLogStore.ts | 21 ++++++++++++------- 1 file changed, 13 insertions(+), 8 deletions(-) diff --git a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts index 73c4a586684..35b406043fb 100644 --- a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts +++ b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts @@ -123,6 +123,7 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { try { await this.storage.write(scope, key, encoded, { atomic: true }); state.storageFailure = undefined; + return true; } catch (error) { state.storageFailure = { error }; throw error; @@ -222,7 +223,7 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { scope: string, key: string, state: LogState, - operation: Promise, + operation: Promise, ): Promise { let owned!: Promise; owned = this.finishOwnedFlush(scope, key, state, operation, () => owned); @@ -234,12 +235,13 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { scope: string, key: string, state: LogState, - operation: Promise, + operation: Promise, owner: () => Promise, ): Promise { let failure: { readonly error: unknown } | undefined; + let wrote = false; try { - await operation; + wrote = await operation; } catch (error) { failure = { error }; } @@ -248,7 +250,7 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { try { if (failure === undefined) { while (state.flushPromise === owned && state.pending.length > 0) { - await this.drain(scope, key, state); + if (await this.drain(scope, key, state)) wrote = true; } } } finally { @@ -258,13 +260,14 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { } } if (failure !== undefined) throw failure.error; - this.writeEmitter.fire({ scope, key }); + if (wrote) this.writeEmitter.fire({ scope, key }); } - private async drain(scope: string, key: string, state: LogState): Promise { + private async drain(scope: string, key: string, state: LogState): Promise { const cutoverEpoch = state.cutoverEpoch; await state.ready; - if (state.cutoverEpoch !== cutoverEpoch) return; + if (state.cutoverEpoch !== cutoverEpoch) return false; + let wrote = false; while (state.pending.length > 0) { const batch = state.pending.slice(); try { @@ -273,9 +276,11 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { const failure = (state.storageFailure ??= { error }); throw failure.error; } - if (state.cutoverEpoch !== cutoverEpoch) return; + wrote = true; + if (state.cutoverEpoch !== cutoverEpoch) return wrote; state.pending.splice(0, batch.length); } + return wrote; } } From 872a7c47067bcdb5f94489846c61210b449a2ca6 Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 15:44:51 +0800 Subject: [PATCH 07/10] fix(agent-core-v2): retry shared snapshot loads that span a write Callers joining an in-flight single-flight load after a completed write still received the pre-write snapshot: the generation check only guarded cache publication, not the value returned to awaiters. Each load now carries the generation it started at, and catalog()/sessionIndex() re-read (coalesced through the same single-flight) when the settled load's generation is stale. --- .../workspaceAliasesService.ts | 32 ++++++---- .../workspaceAliasesService.test.ts | 59 ++++++++++++++++++- 2 files changed, 78 insertions(+), 13 deletions(-) diff --git a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts index 6e08d47bec8..9807ca924eb 100644 --- a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts +++ b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts @@ -48,8 +48,12 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli private catalogCache: CatalogSnapshot | undefined; private sessionIndexCache: SessionIndexSnapshot | undefined; - private catalogPromise: Promise | undefined; - private sessionIndexPromise: Promise | undefined; + private catalogPromise: + | Promise<{ snapshot: CatalogSnapshot; generation: number }> + | undefined; + private sessionIndexPromise: + | Promise<{ snapshot: SessionIndexSnapshot; generation: number }> + | undefined; private invalidationGeneration = 0; private catalogMergePrimed = false; @@ -102,13 +106,15 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli return merged; } - private catalog(): Promise { - if (this.catalogCache !== undefined) return Promise.resolve(this.catalogCache); + private async catalog(): Promise { + if (this.catalogCache !== undefined) return this.catalogCache; this.catalogPromise ??= this.loadCatalog(); - return this.catalogPromise; + const { snapshot, generation } = await this.catalogPromise; + if (generation !== this.invalidationGeneration) return this.catalog(); + return snapshot; } - private async loadCatalog(): Promise { + private async loadCatalog(): Promise<{ snapshot: CatalogSnapshot; generation: number }> { try { if (!this.catalogMergePrimed) { await this.workspaces.list(); @@ -127,19 +133,21 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli if (generation === this.invalidationGeneration) { this.catalogCache = snapshot; } - return snapshot; + return { snapshot, generation }; } finally { this.catalogPromise = undefined; } } - private sessionIndex(): Promise { - if (this.sessionIndexCache !== undefined) return Promise.resolve(this.sessionIndexCache); + private async sessionIndex(): Promise { + if (this.sessionIndexCache !== undefined) return this.sessionIndexCache; this.sessionIndexPromise ??= this.loadSessionIndex(); - return this.sessionIndexPromise; + const { snapshot, generation } = await this.sessionIndexPromise; + if (generation !== this.invalidationGeneration) return this.sessionIndex(); + return snapshot; } - private async loadSessionIndex(): Promise { + private async loadSessionIndex(): Promise<{ snapshot: SessionIndexSnapshot; generation: number }> { try { const generation = this.invalidationGeneration; const entries = await readSessionIndexEntries(this.storage); @@ -151,7 +159,7 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli if (generation === this.invalidationGeneration) { this.sessionIndexCache = snapshot; } - return snapshot; + return { snapshot, generation }; } finally { this.sessionIndexPromise = undefined; } diff --git a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts index 536c9b6a580..3e119c7592c 100644 --- a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts +++ b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts @@ -20,7 +20,7 @@ import { IAppendLogStore } from '#/persistence/interface/appendLogStore'; import { IAtomicDocumentStore } from '#/persistence/interface/atomicDocumentStore'; import { IFileSystemStorageService } from '#/persistence/interface/storage'; import { IEventService } from '#/app/event/event'; -import { IWorkspaceService } from '#/app/workspace/workspace'; +import { IWorkspaceService, type Workspace } from '#/app/workspace/workspace'; import { WorkspaceService } from '#/app/workspace/workspaceService'; import { FileWorkspacePersistence } from '#/app/workspace/fileWorkspacePersistence'; import { @@ -255,6 +255,63 @@ describe('WorkspaceAliasesService (file-backed)', () => { ); }); + it('resolveAliasIds retries a shared load that spanned a write', async () => { + class GatedStorage extends FileStorageService { + catalogReads = 0; + gate: Promise | undefined; + override async read(scope: string, key: string): Promise { + if (key === 'workspaces.json') { + this.catalogReads += 1; + if (this.gate !== undefined) await this.gate; + } + return super.read(scope, key); + } + } + const entry = (root: string): PersistedWorkspaceEntry => ({ + root, + name: 'proj', + created_at: '2026-01-01T00:00:00.000Z', + last_opened_at: '2026-01-01T00:00:00.000Z', + }); + const typedRoot = 'C:\\Users\\Foo\\Proj'; + const typedId = encodeWorkDirKey(typedRoot); + const legacyId = 'wd_proj_deadbeef0002'; + await writeWorkspacesJson({ [typedId]: entry(typedRoot) }); + const storage = new GatedStorage(homeDir); + const aliases = build(undefined, storage); + const persistence = currentHost!.app.accessor.get(IWorkspacePersistence); + const ws = (id: string, root: string): Workspace => ({ + id, + root, + name: 'proj', + createdAt: 0, + lastOpenedAt: 0, + }); + + await aliases.resolveAliasIds(typedId); + await persistence.save({ workspaces: [ws(typedId, typedRoot)], deletedIds: [] }); + + let release: (() => void) | undefined; + storage.gate = new Promise((resolve) => { + release = resolve; + }); + const baseline = storage.catalogReads; + const p1 = aliases.resolveAliasIds(typedId); + await vi.waitFor(() => { + expect(storage.catalogReads).toBe(baseline + 1); + }); + await persistence.save({ + workspaces: [ws(typedId, typedRoot), ws(legacyId, 'c:\\users\\foo\\proj')], + deletedIds: [], + }); + const p2 = aliases.resolveAliasIds(legacyId); + release!(); + + const [r1, r2] = await Promise.all([p1, p2]); + expect(r1.toSorted()).toEqual([legacyId, typedId].toSorted()); + expect(r2.toSorted()).toEqual([legacyId, typedId].toSorted()); + }); + it('resolveAliasIds follows in-process session-index writes synchronously', async () => { const entry = (root: string): PersistedWorkspaceEntry => ({ root, From 8e1f2414477bd819170231312fb9a8228513dc8b Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 16:19:31 +0800 Subject: [PATCH 08/10] fix(agent-core-v2): report partial progress when an append-log drain fails A drain that persisted one batch and then failed the next threw without recording the durable write, so onDidWrite never fired for records that were in fact persisted (the alias session-index snapshot then missed its synchronous invalidation). The write box now threads through the whole owned flush: each successful batch marks it, and the event fires before the failure propagates. --- .../backends/node-fs/appendLogStore.ts | 49 ++++++++++--------- .../backends/node-fs/appendLogStore.test.ts | 31 ++++++++++++ 2 files changed, 58 insertions(+), 22 deletions(-) diff --git a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts index 35b406043fb..e38c7aeec58 100644 --- a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts +++ b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts @@ -129,7 +129,7 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { throw error; } }); - await this.ownFlush(scope, key, state, rewrite); + await this.ownFlush(scope, key, state, rewrite, { value: false }); } async flush(): Promise { @@ -197,7 +197,8 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { private flushState(scope: string, key: string, state: LogState): Promise { if (state.flushPromise !== undefined) return state.flushPromise; if (state.storageFailure !== undefined) return Promise.reject(state.storageFailure.error); - return this.ownFlush(scope, key, state, this.drain(scope, key, state)); + const wroteBox = { value: false }; + return this.ownFlush(scope, key, state, this.drain(scope, key, state, wroteBox), wroteBox); } private release(scope: string, key: string, state: LogState): void { @@ -224,9 +225,10 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { key: string, state: LogState, operation: Promise, + wroteBox: { value: boolean }, ): Promise { let owned!: Promise; - owned = this.finishOwnedFlush(scope, key, state, operation, () => owned); + owned = this.finishOwnedFlush(scope, key, state, operation, wroteBox, () => owned); state.flushPromise = owned; return owned; } @@ -236,34 +238,36 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { key: string, state: LogState, operation: Promise, + wroteBox: { value: boolean }, owner: () => Promise, ): Promise { let failure: { readonly error: unknown } | undefined; - let wrote = false; try { - wrote = await operation; - } catch (error) { - failure = { error }; - } - const owned = owner(); - if (state.flushPromise === owned) { - try { - if (failure === undefined) { - while (state.flushPromise === owned && state.pending.length > 0) { - if (await this.drain(scope, key, state)) wrote = true; - } - } - } finally { - if (state.flushPromise === owned) { - state.flushPromise = undefined; + if (await operation) wroteBox.value = true; + const owned = owner(); + if (state.flushPromise === owned) { + while (state.flushPromise === owned && state.pending.length > 0) { + await this.drain(scope, key, state, wroteBox); } } + } catch (error) { + failure ??= { error }; + } finally { + const owned = owner(); + if (state.flushPromise === owned) { + state.flushPromise = undefined; + } } + if (wroteBox.value) this.writeEmitter.fire({ scope, key }); if (failure !== undefined) throw failure.error; - if (wrote) this.writeEmitter.fire({ scope, key }); } - private async drain(scope: string, key: string, state: LogState): Promise { + private async drain( + scope: string, + key: string, + state: LogState, + wroteBox?: { value: boolean }, + ): Promise { const cutoverEpoch = state.cutoverEpoch; await state.ready; if (state.cutoverEpoch !== cutoverEpoch) return false; @@ -272,11 +276,12 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { const batch = state.pending.slice(); try { await this.storage.append(scope, key, encodeBatch(batch), { durable: true }); + wrote = true; + if (wroteBox !== undefined) wroteBox.value = true; } catch (error) { const failure = (state.storageFailure ??= { error }); throw failure.error; } - wrote = true; if (state.cutoverEpoch !== cutoverEpoch) return wrote; state.pending.splice(0, batch.length); } diff --git a/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts b/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts index 6652fde8997..10be8db29fc 100644 --- a/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts +++ b/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts @@ -759,4 +759,35 @@ describe('AppendLogStore', () => { { n: 2, s: '日本語' }, ]); }); + + it('does not emit onDidWrite when a flush persists nothing', async () => { + const events: string[] = []; + record.onDidWrite((write) => events.push(`${write.scope}/${write.key}`)); + + record.append(SCOPE, KEY, { n: 1 }); + await record.flush(); + expect(events).toEqual([`${SCOPE}/${KEY}`]); + + await record.flush(); + await record.flush(); + expect(events).toEqual([`${SCOPE}/${KEY}`]); + }); + + it('emits onDidWrite for batches persisted before a later drain failure', async () => { + const events: string[] = []; + record.onDidWrite((write) => events.push(`${write.scope}/${write.key}`)); + const original = storage.append.bind(storage); + let calls = 0; + storage.append = async (scope, key, data, options) => { + calls += 1; + if (calls > 1) throw new Error('disk full'); + const result = await original(scope, key, data, options); + record.append(scope, key, { n: 2 }); + return result; + }; + + record.append(SCOPE, KEY, { n: 1 }); + await expect(record.flush()).rejects.toThrow('disk full'); + expect(events).toEqual([`${SCOPE}/${KEY}`]); + }); }); From f8b8f01077a15130c7f6635a409746c0d2ba3d5e Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 16:33:37 +0800 Subject: [PATCH 09/10] fix(agent-core-v2): retry the whole alias resolution across a mid-write The per-snapshot retry guarded each read on its own, so a write landing between the catalog and session-index reads returned an alias set assembled across two generations. resolveAliasIds now captures the generation once, reads both snapshots together, and retries the whole resolution when either input was invalidated mid-flight. The spanned-write test is reworked to gate after the load (so the snapshot content genuinely predates the write), and a new case covers the cross-generation mix directly. --- .../workspaceAliasesService.ts | 29 +++--- .../workspaceAliasesService.test.ts | 94 ++++++++++++++++--- 2 files changed, 99 insertions(+), 24 deletions(-) diff --git a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts index 9807ca924eb..f183bd28b6e 100644 --- a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts +++ b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts @@ -90,20 +90,23 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli } async resolveAliasIds(id: string): Promise { - const catalog = await this.catalog(); - const entry = catalog.byId.get(id); - if (entry === undefined) return [id]; - const rootKey = workspaceRootKey(entry.root); - const index = await this.sessionIndex(); - const fromCatalog = catalog.idsByRootKey.get(rootKey); - const fromIndex = index.idsByRootKey.get(rootKey); - if (fromCatalog === undefined) return fromIndex ?? [id]; - if (fromIndex === undefined) return fromCatalog; - const merged = [...fromCatalog]; - for (const alias of fromIndex) { - if (!merged.includes(alias)) merged.push(alias); + for (;;) { + const generation = this.invalidationGeneration; + const [catalog, index] = await Promise.all([this.catalog(), this.sessionIndex()]); + if (generation !== this.invalidationGeneration) continue; + const entry = catalog.byId.get(id); + if (entry === undefined) return [id]; + const rootKey = workspaceRootKey(entry.root); + const fromCatalog = catalog.idsByRootKey.get(rootKey); + const fromIndex = index.idsByRootKey.get(rootKey); + if (fromCatalog === undefined) return fromIndex ?? [id]; + if (fromIndex === undefined) return fromCatalog; + const merged = [...fromCatalog]; + for (const alias of fromIndex) { + if (!merged.includes(alias)) merged.push(alias); + } + return merged; } - return merged; } private async catalog(): Promise { diff --git a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts index 3e119c7592c..08233f3930b 100644 --- a/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts +++ b/packages/agent-core-v2/test/app/workspaceAliases/workspaceAliasesService.test.ts @@ -26,6 +26,7 @@ import { FileWorkspacePersistence } from '#/app/workspace/fileWorkspacePersisten import { IWorkspacePersistence, type PersistedWorkspaceEntry, + type WorkspaceCatalog, } from '#/app/workspace/workspacePersistence'; import { IWorkspaceAliases } from '#/app/workspaceAliases/workspaceAliases'; import { WorkspaceAliasesService } from '#/app/workspaceAliases/workspaceAliasesService'; @@ -83,11 +84,13 @@ describe('WorkspaceAliasesService (file-backed)', () => { function build( hostFs: IHostFileSystem = new HostFileSystem(), fileStorage: FileStorageService = new FileStorageService(homeDir), + persistence?: IWorkspacePersistence, ): IWorkspaceAliases { const host = createScopedTestHost([ stubPair(IFileSystemStorageService, fileStorage), stubPair(IAtomicDocumentStore, new JsonAtomicDocumentStore(fileStorage)), stubPair(IAppendLogStore, new AppendLogStore(fileStorage)), + ...(persistence !== undefined ? [stubPair(IWorkspacePersistence, persistence)] : []), stubPair(IHostFileSystem, hostFs), stubPair(IEventService, { publish: () => {}, @@ -256,12 +259,78 @@ describe('WorkspaceAliasesService (file-backed)', () => { }); it('resolveAliasIds retries a shared load that spanned a write', async () => { + class GatedPersistence implements IWorkspacePersistence { + declare readonly _serviceBrand: undefined; + loads = 0; + gate: Promise | undefined; + constructor(private readonly inner: IWorkspacePersistence) {} + get onDidChange(): IWorkspacePersistence['onDidChange'] { + return this.inner.onDidChange; + } + async load(): ReturnType { + this.loads += 1; + const catalog = await this.inner.load(); + if (this.gate !== undefined) await this.gate; + return catalog; + } + save(catalog: WorkspaceCatalog): Promise { + return this.inner.save(catalog); + } + } + const entry = (root: string): PersistedWorkspaceEntry => ({ + root, + name: 'proj', + created_at: '2026-01-01T00:00:00.000Z', + last_opened_at: '2026-01-01T00:00:00.000Z', + }); + const typedRoot = 'C:\\Users\\Foo\\Proj'; + const typedId = encodeWorkDirKey(typedRoot); + const legacyId = 'wd_proj_deadbeef0002'; + await writeWorkspacesJson({ [typedId]: entry(typedRoot) }); + const storage = new FileStorageService(homeDir); + const persistence = new GatedPersistence( + new FileWorkspacePersistence(new JsonAtomicDocumentStore(storage)), + ); + const aliases = build(undefined, storage, persistence); + const ws = (id: string, root: string): Workspace => ({ + id, + root, + name: 'proj', + createdAt: 0, + lastOpenedAt: 0, + }); + + await aliases.resolveAliasIds(typedId); + await persistence.save({ workspaces: [ws(typedId, typedRoot)], deletedIds: [] }); + + let release: (() => void) | undefined; + persistence.gate = new Promise((resolve) => { + release = resolve; + }); + const baseline = persistence.loads; + const p1 = aliases.resolveAliasIds(typedId); + await vi.waitFor(() => { + expect(persistence.loads).toBe(baseline + 1); + }); + await persistence.save({ + workspaces: [ws(typedId, typedRoot), ws(legacyId, 'c:\\users\\foo\\proj')], + deletedIds: [], + }); + const p2 = aliases.resolveAliasIds(legacyId); + release!(); + + const [r1, r2] = await Promise.all([p1, p2]); + expect(r1.toSorted()).toEqual([legacyId, typedId].toSorted()); + expect(r2.toSorted()).toEqual([legacyId, typedId].toSorted()); + }); + + it('resolveAliasIds does not mix snapshots across a mid-resolution write', async () => { class GatedStorage extends FileStorageService { - catalogReads = 0; + indexReads = 0; gate: Promise | undefined; override async read(scope: string, key: string): Promise { - if (key === 'workspaces.json') { - this.catalogReads += 1; + if (key === 'session_index.jsonl') { + this.indexReads += 1; if (this.gate !== undefined) await this.gate; } return super.read(scope, key); @@ -280,6 +349,7 @@ describe('WorkspaceAliasesService (file-backed)', () => { const storage = new GatedStorage(homeDir); const aliases = build(undefined, storage); const persistence = currentHost!.app.accessor.get(IWorkspacePersistence); + const appendLogs = currentHost!.app.accessor.get(IAppendLogStore); const ws = (id: string, root: string): Workspace => ({ id, root, @@ -289,27 +359,29 @@ describe('WorkspaceAliasesService (file-backed)', () => { }); await aliases.resolveAliasIds(typedId); - await persistence.save({ workspaces: [ws(typedId, typedRoot)], deletedIds: [] }); + appendLogs.append('', 'session_index.jsonl', { + sessionId: 's9', + sessionDir: 'sessions/s/s9', + workDir: join(homeDir, 'unrelated'), + }); + await appendLogs.flush(); let release: (() => void) | undefined; storage.gate = new Promise((resolve) => { release = resolve; }); - const baseline = storage.catalogReads; - const p1 = aliases.resolveAliasIds(typedId); + const baseline = storage.indexReads; + const p = aliases.resolveAliasIds(typedId); await vi.waitFor(() => { - expect(storage.catalogReads).toBe(baseline + 1); + expect(storage.indexReads).toBe(baseline + 1); }); await persistence.save({ workspaces: [ws(typedId, typedRoot), ws(legacyId, 'c:\\users\\foo\\proj')], deletedIds: [], }); - const p2 = aliases.resolveAliasIds(legacyId); release!(); - const [r1, r2] = await Promise.all([p1, p2]); - expect(r1.toSorted()).toEqual([legacyId, typedId].toSorted()); - expect(r2.toSorted()).toEqual([legacyId, typedId].toSorted()); + expect((await p).toSorted()).toEqual([legacyId, typedId].toSorted()); }); it('resolveAliasIds follows in-process session-index writes synchronously', async () => { From 3411fd134b841d8a4ea9889e8cb7cb4663c881b1 Mon Sep 17 00:00:00 2001 From: liruifengv Date: Fri, 28 Aug 2026 17:10:06 +0800 Subject: [PATCH 10/10] fix(agent-core-v2): replace the session-index fs watch with a size check A resident chokidar watcher per server on the shared home directory degraded watch delivery for unrelated files under test-suite boot volume (the prompts suite lost the config.toml reload race and the catalog missed a just-written model). In-process appends were already covered synchronously by the append-log write event; cross-process writers now surface through a per-call size comparison on the append-only file, which costs one stat per resolve and needs no resident watcher. --- .../workspaceAliasesService.ts | 24 +++++++++---------- 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts index f183bd28b6e..d9ceb425fb9 100644 --- a/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts +++ b/packages/agent-core-v2/src/app/workspaceAliases/workspaceAliasesService.ts @@ -47,7 +47,7 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli declare readonly _serviceBrand: undefined; private catalogCache: CatalogSnapshot | undefined; - private sessionIndexCache: SessionIndexSnapshot | undefined; + private sessionIndexCache: { snapshot: SessionIndexSnapshot; size: number | undefined } | undefined; private catalogPromise: | Promise<{ snapshot: CatalogSnapshot; generation: number }> | undefined; @@ -78,15 +78,6 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli } }), ); - const watchSessionIndex = this.storage.watch?.(SESSION_INDEX_SCOPE, SESSION_INDEX_KEY); - if (watchSessionIndex !== undefined) { - this._register( - watchSessionIndex(() => { - this.invalidationGeneration += 1; - this.sessionIndexCache = undefined; - }), - ); - } } async resolveAliasIds(id: string): Promise { @@ -143,7 +134,13 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli } private async sessionIndex(): Promise { - if (this.sessionIndexCache !== undefined) return this.sessionIndexCache; + const cache = this.sessionIndexCache; + if ( + cache !== undefined && + (await this.storage.size(SESSION_INDEX_SCOPE, SESSION_INDEX_KEY)) === cache.size + ) { + return cache.snapshot; + } this.sessionIndexPromise ??= this.loadSessionIndex(); const { snapshot, generation } = await this.sessionIndexPromise; if (generation !== this.invalidationGeneration) return this.sessionIndex(); @@ -160,7 +157,10 @@ export class WorkspaceAliasesService extends Disposable implements IWorkspaceAli ), }; if (generation === this.invalidationGeneration) { - this.sessionIndexCache = snapshot; + this.sessionIndexCache = { + snapshot, + size: await this.storage.size(SESSION_INDEX_SCOPE, SESSION_INDEX_KEY), + }; } return { snapshot, generation }; } finally {