diff --git a/apps/desktop/src/main/cindy-brain/__tests__/ghostOauthFlow.test.ts b/apps/desktop/src/main/cindy-brain/__tests__/ghostOauthFlow.test.ts index 211271928d..c2e7b38780 100644 --- a/apps/desktop/src/main/cindy-brain/__tests__/ghostOauthFlow.test.ts +++ b/apps/desktop/src/main/cindy-brain/__tests__/ghostOauthFlow.test.ts @@ -33,15 +33,25 @@ function jsonResponse(body: unknown, status = 200): Response { } /** 从 openExternal 捕获的授权 URL 中提取回调地址与 state,模拟浏览器完成授权。 */ -function browserRedirect(authorizeUrl: string, params: (u: URL) => Record): void { +function browserRedirect( + authorizeUrl: string, + params: (u: URL) => Record, + callbackFetch: typeof fetch = fetch, +): Promise { const url = new URL(authorizeUrl); const redirectUri = url.searchParams.get('redirect_uri'); if (!redirectUri) throw new Error('authorize URL 缺 redirect_uri'); const cb = new URL(redirectUri); for (const [k, v] of Object.entries(params(url))) cb.searchParams.set(k, v); - // 不 await:引擎在 race 回调,fire-and-forget 即可;失败让测试超时暴露。 - setImmediate(() => { - void fetch(cb.toString()).catch(() => undefined); + // openExternal 会 await 本 Promise:loopback 传输失败必须原样暴露,不能吞掉后 + // 让 OAuth flow 等到 Vitest 的外层超时。消费 body 后再结算,避免连接悬空。 + return new Promise((resolve, reject) => { + setImmediate(() => { + void (async () => { + const response = await callbackFetch(cb.toString()); + await response.arrayBuffer(); + })().then(resolve, reject); + }); }); } @@ -93,6 +103,19 @@ function isListenFailed(value: unknown): boolean { } describe('startGhostOauthFlow', () => { + it('browserRedirect 把 loopback callback 传输失败直接交给调用方', async () => { + const authorizeUrl = new URL(BASE_CONFIG.authorizeUrl); + authorizeUrl.searchParams.set('redirect_uri', 'http://127.0.0.1:1/callback'); + const callbackFailure = new Error('loopback callback unavailable'); + const callbackFetch = vi.fn(async () => { + throw callbackFailure; + }) as unknown as typeof fetch; + + await expect( + browserRedirect(authorizeUrl.toString(), () => ({ code: 'c0', state: 's0' }), callbackFetch), + ).rejects.toBe(callbackFailure); + }); + it('happy path:PKCE + state 校验 + code 换 token', async () => { let capturedAuthorizeUrl = ''; const fetchImpl = vi.fn(async (input: string | URL | Request, init?: RequestInit) => { @@ -130,7 +153,7 @@ describe('startGhostOauthFlow', () => { expect(u.searchParams.get('response_type')).toBe('code'); expect(u.searchParams.get('scope')).toBe('read:a write:b'); expect(u.searchParams.get('code_challenge_method')).toBe('S256'); - browserRedirect(url, (au) => ({ + return browserRedirect(url, (au) => ({ code: 'code-abc', state: au.searchParams.get('state') ?? '', })); @@ -174,7 +197,7 @@ describe('startGhostOauthFlow', () => { // 保留参数不被意识声明顶掉。 expect(u.searchParams.get('client_id')).toBe('client-123'); expect(u.searchParams.get('state')).not.toBe('EVIL-STATE'); - browserRedirect(url, (au) => ({ + return browserRedirect(url, (au) => ({ code: 'c2', state: au.searchParams.get('state') ?? '', })); @@ -201,7 +224,10 @@ describe('startGhostOauthFlow', () => { fetchImpl: fetchImpl as unknown as typeof fetch, openExternal: (url) => { expect(new URL(url).searchParams.get('code_challenge')).toBeNull(); - browserRedirect(url, (au) => ({ code: 'c3', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url, (au) => ({ + code: 'c3', + state: au.searchParams.get('state') ?? '', + })); }, }); expect(result.ok).toBe(true); @@ -246,7 +272,7 @@ describe('startGhostOauthFlow', () => { fetchImpl: vi.fn() as unknown as typeof fetch, openExternal: (url) => { // state 先行校验后,error 参数只在 state 匹配时才结算(与 grok 同口径)。 - browserRedirect(url, (au) => ({ + return browserRedirect(url, (au) => ({ error: 'access_denied', state: au.searchParams.get('state') ?? '', })); @@ -314,7 +340,10 @@ describe('startGhostOauthFlow', () => { config: BASE_CONFIG, fetchImpl: vi.fn(async () => jsonResponse({ access_token: 'at-6' })) as unknown as typeof fetch, openExternal: (url2) => { - browserRedirect(url2, (au) => ({ code: 'c6', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url2, (au) => ({ + code: 'c6', + state: au.searchParams.get('state') ?? '', + })); }, }); }, @@ -330,7 +359,10 @@ describe('startGhostOauthFlow', () => { jsonResponse({ error: 'invalid_client', error_description: 'bad client' }, 401), ) as unknown as typeof fetch, openExternal: (url) => { - browserRedirect(url, (au) => ({ code: 'c7', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url, (au) => ({ + code: 'c7', + state: au.searchParams.get('state') ?? '', + })); }, }); expect(result).toMatchObject({ ok: false, error: 'EXCHANGE_FAILED' }); @@ -348,7 +380,10 @@ describe('startGhostOauthFlow', () => { openExternal: (url) => { const u = new URL(url); expect(u.searchParams.get('redirect_uri')).toBe(`http://127.0.0.1:${freePort}/callback`); - browserRedirect(url, (au) => ({ code: 'c-fixed', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url, (au) => ({ + code: 'c-fixed', + state: au.searchParams.get('state') ?? '', + })); }, }), ]); @@ -395,7 +430,10 @@ describe('startGhostOauthFlow', () => { config: { ...BASE_CONFIG, pkce: false, redirectPort: fixedPort }, fetchImpl: vi.fn(async () => jsonResponse({ access_token: 'at-heal' })) as unknown as typeof fetch, openExternal: (url2) => { - browserRedirect(url2, (au) => ({ code: 'c-heal', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url2, (au) => ({ + code: 'c-heal', + state: au.searchParams.get('state') ?? '', + })); }, }); }, @@ -429,7 +467,10 @@ describe('startGhostOauthFlow', () => { config: { ...BASE_CONFIG, pkce: false, redirectPort: fixedPort }, fetchImpl: vi.fn(async () => jsonResponse({ access_token: 'at-third' })) as unknown as typeof fetch, openExternal: (url3) => { - browserRedirect(url3, (au) => ({ code: 'c-third', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url3, (au) => ({ + code: 'c-third', + state: au.searchParams.get('state') ?? '', + })); }, }); }, @@ -461,7 +502,10 @@ describe('startGhostOauthFlow', () => { config: { ...BASE_CONFIG, pkce: false, redirectPort: heldPort }, fetchImpl: vi.fn(async () => jsonResponse({ access_token: 'at-reclaim' })) as unknown as typeof fetch, openExternal: (url) => { - browserRedirect(url, (au) => ({ code: 'c-reclaim', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url, (au) => ({ + code: 'c-reclaim', + state: au.searchParams.get('state') ?? '', + })); }, reclaimPort, }); @@ -528,7 +572,10 @@ describe('startGhostOauthFlow', () => { openExternal: (url) => { challengeFromUrl = new URL(url).searchParams.get('code_challenge'); expect(challengeFromUrl).toBeTruthy(); - browserRedirect(url, (au) => ({ code: 'c-broker', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url, (au) => ({ + code: 'c-broker', + state: au.searchParams.get('state') ?? '', + })); }, }); expect(result).toMatchObject({ ok: true }); @@ -551,7 +598,7 @@ describe('startGhostOauthFlow', () => { fetchImpl: vi.fn() as unknown as typeof fetch, broker, openExternal: (url) => { - browserRedirect(url, (authorizeUrl) => ({ + return browserRedirect(url, (authorizeUrl) => ({ code: 'c-unavailable', state: authorizeUrl.searchParams.get('state') ?? '', })); @@ -580,7 +627,10 @@ describe('startGhostOauthFlow', () => { broker, openExternal: (url) => { expect(new URL(url).searchParams.get('code_challenge')).toBeNull(); - browserRedirect(url, (au) => ({ code: 'c-np', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url, (au) => ({ + code: 'c-np', + state: au.searchParams.get('state') ?? '', + })); }, }); expect(result).toMatchObject({ ok: true }); @@ -826,7 +876,10 @@ describe('startGhostOauthFlow', () => { fetchImpl: vi.fn(async () => jsonResponse({ access_token: 'at-sd' })) as unknown as typeof fetch, openExternal: (url) => { expect(new URL(url).searchParams.get('scope')).toBe('read:a,write:b'); - browserRedirect(url, (au) => ({ code: 'c-sd', state: au.searchParams.get('state') ?? '' })); + return browserRedirect(url, (au) => ({ + code: 'c-sd', + state: au.searchParams.get('state') ?? '', + })); }, }); expect(result).toMatchObject({ ok: true }); diff --git a/apps/desktop/src/main/maker-host/__tests__/pi-package-store-security.test.ts b/apps/desktop/src/main/maker-host/__tests__/pi-package-store-security.test.ts index e9911a9f96..4b6a60e299 100644 --- a/apps/desktop/src/main/maker-host/__tests__/pi-package-store-security.test.ts +++ b/apps/desktop/src/main/maker-host/__tests__/pi-package-store-security.test.ts @@ -8,6 +8,7 @@ import { import { createRequire } from 'node:module'; import os from 'node:os'; import path from 'node:path'; +import type { PiPackageMutationRequest } from '../../../shared/piPackages.js'; // Windows 未开启 Developer Mode / 无 Create Symbolic Link 权限时,文件 symlink // 会返回 EPERM;目录 junction 不受此限制,但本文件的竞态用例必须替换单个文件。 @@ -51,6 +52,10 @@ const lockRuntime = vi.hoisted(() => ({ nextStatus: null as null | { held: false; reason: 'busy' | 'unavailable' }, })); +const loggerRuntime = vi.hoisted(() => ({ + info: vi.fn(), +})); + vi.mock('electron', () => ({ app: { getPath: () => runtime.userData }, })); @@ -61,7 +66,7 @@ vi.mock('../../agent-binaries/index.js', () => ({ vi.mock('../../logger.js', () => ({ createLogger: () => ({ - trace: vi.fn(), debug: vi.fn(), info: vi.fn(), warn: vi.fn(), error: vi.fn(), fatal: vi.fn(), + trace: vi.fn(), debug: vi.fn(), info: loggerRuntime.info, warn: vi.fn(), error: vi.fn(), fatal: vi.fn(), child() { return this; }, }), })); @@ -212,6 +217,7 @@ beforeEach(async () => { lockRuntime.active = 0; lockRuntime.maxActive = 0; lockRuntime.nextStatus = null; + loggerRuntime.info.mockReset(); vi.resetModules(); }); @@ -1345,7 +1351,7 @@ describe('Pi package executable-code boundary', () => { await expect(store.listPiPackages()).resolves.toMatchObject({ packages: [ { source: sources[0], enabled: true }, - { source: sources[1], enabled: false, warning: 'inspection-limit' }, + { source: sources[1], enabled: true }, ], }); }); @@ -1407,16 +1413,14 @@ describe('Pi package executable-code boundary', () => { await expect(store.listPiPackages()).resolves.toMatchObject({ packages: [ { source: first.source, enabled: true }, - { source: second.source, enabled: false, warning: 'inspection-limit' }, + { source: second.source, enabled: true }, ], }); const state = JSON.parse(await fs.readFile( path.join(runtime.userData, 'pi-package-home', 'cindy-package-state.json'), 'utf8', - )) as { snapshotUnavailableRoots: Record }; - expect(state.snapshotUnavailableRoots).toEqual({ - [await fs.realpath(second.root)]: 'inspection-limit', - }); + )) as { snapshotUnavailablePackages: unknown[] }; + expect(state.snapshotUnavailablePackages).toEqual([]); }); it('isolates current and later roots when aggregate fingerprint budget is exhausted', async () => { @@ -1478,12 +1482,22 @@ describe('Pi package executable-code boundary', () => { await expect(store.listPiPackages()).resolves.toMatchObject({ packages: [ { source: packages[0]!.source, enabled: true }, - { source: packages[1]!.source, enabled: false, warning: 'inspection-limit' }, - { source: packages[2]!.source, enabled: false, warning: 'inspection-limit' }, + { source: packages[1]!.source, enabled: true }, + { source: packages[2]!.source, enabled: true }, ], }); }); + it.each(['entries', 'bytes'] as const)( + 'does not classify an aggregate fingerprint %s limit as a durable package failure', + async (reason) => { + const store = await import('../pi-package-store.js'); + + expect(store.__testing.isDurableSnapshotLimit('aggregate', reason)).toBe(false); + expect(store.__testing.isDurableSnapshotLimit('package', reason)).toBe(true); + }, + ); + it('omits resources owned by a skipped descendant instead of mapping them through a copied ancestor', async () => { const ancestorRoot = await fs.mkdtemp(path.join(os.tmpdir(), 'cindy-pi-package-budget-overlap-')); roots.push(ancestorRoot); @@ -1574,17 +1588,525 @@ describe('Pi package executable-code boundary', () => { { enabled: false, warning: 'inspection-limit' }, ], }); + const stateFile = await fs.readFile( + path.join(runtime.userData, 'pi-package-home', 'cindy-package-state.json'), + 'utf8', + ).catch((error: NodeJS.ErrnoException) => { + if (error.code === 'ENOENT') return null; + throw error; + }); + if (stateFile) { + const state = JSON.parse(stateFile) as { + disabledSources: string[]; + snapshotUnavailablePackages: unknown[]; + }; + expect(state.disabledSources).toEqual([]); + expect(state.snapshotUnavailablePackages).toEqual([]); + } + + vi.resetModules(); + const nextStore = await import('../pi-package-store.js'); + const recoveredRoot = path.join(runtime.userData, 'aggregate-package-recovered'); + await expect(nextStore.resolveManagedPiPackageResources({ + snapshotRoot: recoveredRoot, + snapshotLimits: { + maxEntries: 100, + maxBytes: 1024 * 1024, + maxDurationMs: 10_000, + }, + })).resolves.toMatchObject({ + skills: [ + expect.objectContaining({ + path: path.join(recoveredRoot, '0', 'skills', 'skill-0', 'SKILL.md'), + }), + expect.objectContaining({ + path: path.join(recoveredRoot, '1', 'skills', 'skill-1', 'SKILL.md'), + }), + expect.objectContaining({ + path: path.join(recoveredRoot, '2', 'skills', 'skill-2', 'SKILL.md'), + }), + ], + packageRoots: [ + path.join(recoveredRoot, '0'), + path.join(recoveredRoot, '1'), + path.join(recoveredRoot, '2'), + ], + }); + }); + + it('reuses a persisted inspection-limit without walking the package tree for every session', async () => { + const { root, source } = await createSkillOnlyPackage('npm:oversized-skill'); + const store = await import('../pi-package-store.js'); + const snapshotLimits = { + maxEntries: 1, + maxBytes: 1024 * 1024, + maxDurationMs: 10_000, + }; + + await expect( + store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'first-oversized-skill-snapshot'), + snapshotLimits, + }), + ).resolves.toEqual({ + extensions: [], + skills: [], + promptTemplates: [], + packageRoots: [], + }); + await expect(store.listPiPackages()).resolves.toMatchObject({ + packages: [{ source, enabled: false, warning: 'inspection-limit' }], + }); + + const canonicalRoot = await fs.realpath(root); + const opendirSpy = vi.spyOn(fs, 'opendir'); + try { + await expect( + store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'second-oversized-skill-snapshot'), + snapshotLimits, + }), + ).resolves.toEqual({ + extensions: [], + skills: [], + promptTemplates: [], + packageRoots: [], + }); + const repeatedPackageWalks = opendirSpy.mock.calls.filter(([candidate]) => { + const resolved = path.resolve(String(candidate)); + return resolved === canonicalRoot || resolved.startsWith(`${canonicalRoot}${path.sep}`); + }); + expect(repeatedPackageWalks).toEqual([]); + } finally { + opendirSpy.mockRestore(); + } + }); + + it('does not reuse a persisted inspection-limit for a different source at the same root', async () => { + const { root } = await createSkillOnlyPackage('npm:oversized-source-a'); + const store = await import('../pi-package-store.js'); + await expect( + store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'source-a-limited-snapshot'), + snapshotLimits: { + maxEntries: 1, + maxBytes: 1024 * 1024, + maxDurationMs: 10_000, + }, + }), + ).resolves.toMatchObject({ packageRoots: [] }); + + const replacementSource = 'npm:replacement-source-b'; + runtime.listOutput = `User packages:\n ${replacementSource}\n ${root}\n`; + vi.resetModules(); + const replacementStore = await import('../pi-package-store.js'); + await expect(replacementStore.listPiPackages()).resolves.toMatchObject({ + packages: [ + expect.objectContaining({ + source: replacementSource, + enabled: true, + }), + ], + }); + const recoveredRoot = path.join(runtime.userData, 'source-b-recovered-snapshot'); + await expect( + replacementStore.resolveManagedPiPackageResources({ + snapshotRoot: recoveredRoot, + }), + ).resolves.toMatchObject({ + skills: [ + expect.objectContaining({ + path: path.join(recoveredRoot, '0', 'skills', 'managed-skill', 'SKILL.md'), + }), + ], + packageRoots: [path.join(recoveredRoot, '0')], + }); + }); + + it('retries a legacy v4 inspection-limit without retry metadata once', async () => { + const { source } = await createSkillOnlyPackage('npm:legacy-v4-snapshot-limit'); + const store = await import('../pi-package-store.js'); + await expect(store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'legacy-v4-limited-snapshot'), + snapshotLimits: { + maxEntries: 1, + maxBytes: 1024 * 1024, + maxDurationMs: 10_000, + }, + })).resolves.toMatchObject({ packageRoots: [] }); + + const statePath = path.join( + runtime.userData, + 'pi-package-home', + 'cindy-package-state.json', + ); + const state = JSON.parse(await fs.readFile(statePath, 'utf8')) as { + snapshotUnavailablePackages: Array<{ retryAfterEpochMs?: number }>; + }; + expect(state.snapshotUnavailablePackages).toHaveLength(1); + delete state.snapshotUnavailablePackages[0]!.retryAfterEpochMs; + await fs.writeFile(statePath, JSON.stringify(state)); + + vi.resetModules(); + const nextStore = await import('../pi-package-store.js'); + const recoveredRoot = path.join(runtime.userData, 'legacy-v4-recovered-snapshot'); + await expect(nextStore.resolveManagedPiPackageResources({ + snapshotRoot: recoveredRoot, + })).resolves.toMatchObject({ + skills: [expect.objectContaining({ + path: path.join(recoveredRoot, '0', 'skills', 'managed-skill', 'SKILL.md'), + })], + packageRoots: [path.join(recoveredRoot, '0')], + }); + await expect(nextStore.listPiPackages()).resolves.toMatchObject({ + packages: [expect.objectContaining({ source, enabled: true })], + }); + }); + + it('persists a deterministic inspection-stage limit before the next session walks metadata', async () => { + const { root, source } = await createPackage({ oversizedManifest: true }); + const store = await import('../pi-package-store.js'); + await expect(store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'inspection-limit-first-snapshot'), + })).resolves.toEqual({ + extensions: [], + skills: [], + promptTemplates: [], + packageRoots: [], + }); + + const statePath = path.join( + runtime.userData, + 'pi-package-home', + 'cindy-package-state.json', + ); + await expect(fs.readFile(statePath, 'utf8')).resolves.toSatisfy((raw) => { + const state = JSON.parse(raw) as { + snapshotUnavailablePackages: Array<{ source: string }>; + }; + return state.snapshotUnavailablePackages.some((entry) => entry.source === source); + }); + + vi.resetModules(); + const manifestPath = path.join(root, 'package.json'); + const canonicalManifestPath = path.resolve(await fs.realpath(manifestPath)); + const openSpy = vi.spyOn(fs, 'open'); + const opendirSpy = vi.spyOn(fs, 'opendir'); + try { + const nextStore = await import('../pi-package-store.js'); + await nextStore.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'inspection-limit-second-snapshot'), + }); + // The durable negative identity includes one bounded package.json digest + // so in-place upgrades invalidate it. Reuse must still avoid walking, + // fingerprinting, or copying the package tree. + const openedPaths = await Promise.all(openSpy.mock.calls.map(async ([candidate]) => ( + path.resolve(await fs.realpath(String(candidate))) + ))); + expect(openedPaths.filter((candidate) => candidate === canonicalManifestPath)).toHaveLength(1); + expect(opendirSpy).not.toHaveBeenCalled(); + } finally { + openSpy.mockRestore(); + opendirSpy.mockRestore(); + } + }); + + it('does not spread an inspection-stage failure to a healthy npm sibling on the shared root', async () => { + const npmRoot = path.join(runtime.userData, 'pi-package-home', 'npm'); + const oversizedRoot = path.join(npmRoot, 'node_modules', 'inspection-oversized'); + const healthyRoot = path.join(npmRoot, 'node_modules', 'inspection-healthy'); + await fs.mkdir(path.join(healthyRoot, 'skills', 'healthy'), { recursive: true }); + await fs.mkdir(oversizedRoot, { recursive: true }); + const oversizedSource = 'npm:inspection-oversized'; + const healthySource = 'npm:inspection-healthy'; + await fs.writeFile(path.join(oversizedRoot, 'package.json'), JSON.stringify({ + name: 'inspection-oversized', + version: '1.0.0', + pi: { + prompts: Array.from({ length: 257 }, (_, index) => `prompts/${index}.md`), + }, + })); + await fs.writeFile(path.join(healthyRoot, 'package.json'), JSON.stringify({ + name: 'inspection-healthy', + version: '1.0.0', + pi: { skills: ['./skills'] }, + })); + await fs.writeFile( + path.join(healthyRoot, 'skills', 'healthy', 'SKILL.md'), + '# Healthy\n', + ); + runtime.listOutput = [ + 'User packages:', + ` ${oversizedSource}`, + ` ${oversizedRoot}`, + ` ${healthySource}`, + ` ${healthyRoot}`, + '', + ].join('\n'); + + const store = await import('../pi-package-store.js'); + await store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'inspection-shared-root-snapshot'), + }); + await expect(store.listPiPackages()).resolves.toMatchObject({ + packages: [ + { source: oversizedSource, enabled: false, warning: 'inspection-limit' }, + { source: healthySource, enabled: true }, + ], + }); const state = JSON.parse(await fs.readFile( path.join(runtime.userData, 'pi-package-home', 'cindy-package-state.json'), 'utf8', - )) as { disabledSources: string[]; snapshotUnavailableRoots: Record }; - expect(state.disabledSources).toEqual([]); - expect(state.snapshotUnavailableRoots).toEqual({ - [await fs.realpath(packageRoots[1]!)]: 'inspection-limit', - [await fs.realpath(packageRoots[2]!)]: 'inspection-limit', + )) as { snapshotUnavailablePackages: Array<{ source: string }> }; + expect(state.snapshotUnavailablePackages.map((entry) => entry.source)).toEqual([ + oversizedSource, + ]); + }); + + it('retries an inspection-stage duration limit instead of persisting it', async () => { + const { root, source } = await createSkillOnlyPackage('npm:inspection-duration-retry'); + const store = await import('../pi-package-store.js'); + const originalReaddir = fs.readdir.bind(fs); + const canonicalSkillsRoot = path.resolve(await fs.realpath(path.join(root, 'skills'))); + let now = Date.now(); + let delayedInspection = false; + const nowSpy = vi.spyOn(Date, 'now').mockImplementation(() => now); + const readdirSpy = vi.spyOn(fs, 'readdir').mockImplementation(async (...args) => { + const result = await originalReaddir(...args as Parameters); + const canonicalCandidate = path.resolve(await fs.realpath(String(args[0]))); + if ( + !delayedInspection + && canonicalCandidate === canonicalSkillsRoot + ) { + delayedInspection = true; + now += 2_001; + } + return result; }); + try { + await expect(store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'inspection-duration-first-snapshot'), + })).resolves.toMatchObject({ packageRoots: [] }); + await expect(store.listPiPackages()).resolves.toMatchObject({ + packages: [{ source, enabled: false, warning: 'inspection-limit' }], + }); + + const stateFile = await fs.readFile( + path.join(runtime.userData, 'pi-package-home', 'cindy-package-state.json'), + 'utf8', + ).catch(() => null); + if (stateFile) { + expect(JSON.parse(stateFile)).toMatchObject({ snapshotUnavailablePackages: [] }); + } + + const recoveredRoot = path.join(runtime.userData, 'inspection-duration-recovered-snapshot'); + await expect(store.resolveManagedPiPackageResources({ + snapshotRoot: recoveredRoot, + })).resolves.toMatchObject({ + skills: [ + expect.objectContaining({ + path: path.join(recoveredRoot, '0', 'skills', 'managed-skill', 'SKILL.md'), + }), + ], + packageRoots: [path.join(recoveredRoot, '0')], + }); + } finally { + readdirSpy.mockRestore(); + nowSpy.mockRestore(); + } }); + it('does not reuse a persisted inspection-limit after the installation is replaced in place', async () => { + const { root, source } = await createSkillOnlyPackage('npm:replaced-installation'); + const store = await import('../pi-package-store.js'); + await expect( + store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'old-installation-limited-snapshot'), + snapshotLimits: { + maxEntries: 1, + maxBytes: 1024 * 1024, + maxDurationMs: 10_000, + }, + }), + ).resolves.toMatchObject({ packageRoots: [] }); + + await fs.rm(root, { recursive: true, force: true }); + await fs.mkdir(path.join(root, 'skills', 'managed-skill'), { recursive: true }); + await fs.writeFile( + path.join(root, 'package.json'), + JSON.stringify({ + name: source.slice(4), + version: '2.0.0', + pi: { skills: ['./skills'] }, + }), + ); + await fs.writeFile( + path.join(root, 'skills', 'managed-skill', 'SKILL.md'), + '# Replaced installation\n', + ); + + vi.resetModules(); + const replacementStore = await import('../pi-package-store.js'); + await expect(replacementStore.listPiPackages()).resolves.toMatchObject({ + packages: [ + expect.objectContaining({ + source, + enabled: true, + }), + ], + }); + const recoveredRoot = path.join(runtime.userData, 'new-installation-recovered-snapshot'); + await expect( + replacementStore.resolveManagedPiPackageResources({ + snapshotRoot: recoveredRoot, + }), + ).resolves.toMatchObject({ + skills: [ + expect.objectContaining({ + path: path.join(recoveredRoot, '0', 'skills', 'managed-skill', 'SKILL.md'), + }), + ], + packageRoots: [path.join(recoveredRoot, '0')], + }); + }); + + it('retries a persisted inspection-limit after deep package content shrinks under a stable root identity', async () => { + const { root, source } = await createSkillOnlyPackage('npm:mutated-installation'); + const deepPayload = path.join(root, 'skills', 'managed-skill', 'payload.bin'); + await fs.writeFile(deepPayload, Buffer.alloc(4_096)); + let now = Date.now(); + const nowSpy = vi.spyOn(Date, 'now').mockImplementation(() => now); + const originalLstat = fs.lstat.bind(fs); + const originalStat = fs.stat.bind(fs); + let lstatSpy: { mockRestore(): void } | undefined; + let statSpy: { mockRestore(): void } | undefined; + try { + const store = await import('../pi-package-store.js'); + await expect( + store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'mutated-installation-limited-snapshot'), + snapshotLimits: { + maxEntries: 100, + maxBytes: 1_024, + maxDurationMs: 10_000, + }, + }), + ).resolves.toMatchObject({ packageRoots: [] }); + + const canonicalRoot = await fs.realpath(root); + const stableRootStat = await fs.stat(canonicalRoot); + await fs.writeFile(deepPayload, Buffer.from('.')); + now += 24 * 60 * 60 * 1_000; + + // Root directory metadata is not a content identity. Pin it to the value + // captured for the failed installation so this regression stays + // deterministic on filesystems that happen to touch directory timestamps. + lstatSpy = vi.spyOn(fs, 'lstat').mockImplementation(async (target, options) => ( + path.resolve(String(target)) === canonicalRoot + ? stableRootStat + : originalLstat(target, options as never) + )); + statSpy = vi.spyOn(fs, 'stat').mockImplementation(async (target, options) => ( + path.resolve(String(target)) === canonicalRoot + ? stableRootStat + : originalStat(target, options as never) + )); + vi.resetModules(); + const replacementStore = await import('../pi-package-store.js'); + const recoveredRoot = path.join(runtime.userData, 'mutated-installation-recovered-snapshot'); + await expect(replacementStore.resolveManagedPiPackageResources({ + snapshotRoot: recoveredRoot, + snapshotLimits: { + maxEntries: 100, + maxBytes: 1_024, + maxDurationMs: 10_000, + }, + })).resolves.toMatchObject({ + skills: [expect.objectContaining({ + path: path.join(recoveredRoot, '0', 'skills', 'managed-skill', 'SKILL.md'), + })], + packageRoots: [path.join(recoveredRoot, '0')], + }); + } finally { + lstatSpy?.mockRestore(); + statSpy?.mockRestore(); + nowSpy.mockRestore(); + } + }); + + it.each([ + ['install', 'install', undefined], + ['update', 'update', undefined], + ['remove', 'remove', undefined], + ['enable', 'set-enabled', true], + ['disable', 'set-enabled', false], + ] as const)( + 'a production %s mutation clears the durable snapshot failure before fresh inspection', + async (_label, action, enabled) => { + const { source } = await createSkillOnlyPackage( + `npm:mutation-recovery-${action}-${String(enabled)}`, + ); + const store = await import('../pi-package-store.js'); + const enableIdentity = await store.capturePiPackageEnableIdentity(source); + await expect( + store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, `limited-before-${action}-${String(enabled)}`), + snapshotLimits: { + maxEntries: 1, + maxBytes: 1024 * 1024, + maxDurationMs: 10_000, + }, + }), + ).resolves.toMatchObject({ packageRoots: [] }); + + const request: PiPackageMutationRequest = + action === 'set-enabled' ? { action, source, enabled: enabled! } : { action, source }; + if (action === 'set-enabled' && enabled === true) { + const { issuePiPackageMutationGrant } = await import('../pi-package-mutation-grant.js'); + await store.mutatePiPackage( + request, + issuePiPackageMutationGrant(request, { + expectedPackageFingerprint: enableIdentity.expectedPackageFingerprint, + }), + ); + } else if (action === 'set-enabled') { + await store.mutatePiPackage(request); + } else { + await mutateAuthorized(store, request); + } + + if (enabled === false) { + await expect(store.listPiPackages()).resolves.toMatchObject({ + packages: [{ source, enabled: false }], + }); + } else { + const recoveredRoot = path.join( + runtime.userData, + `recovered-after-${action}-${String(enabled)}`, + ); + await expect( + store.resolveManagedPiPackageResources({ + snapshotRoot: recoveredRoot, + }), + ).resolves.toMatchObject({ + skills: [ + expect.objectContaining({ + path: path.join(recoveredRoot, '0', 'skills', 'managed-skill', 'SKILL.md'), + }), + ], + packageRoots: [path.join(recoveredRoot, '0')], + }); + } + const state = JSON.parse( + await fs.readFile( + path.join(runtime.userData, 'pi-package-home', 'cindy-package-state.json'), + 'utf8', + ), + ) as { snapshotUnavailablePackages: unknown[] }; + expect(state.snapshotUnavailablePackages).toEqual([]); + }, + ); + it('does not run a skipped descendant extension from an unverified ancestor snapshot', async () => { const ancestorRoot = await fs.mkdtemp(path.join(os.tmpdir(), 'cindy-pi-package-approved-overlap-')); roots.push(ancestorRoot); @@ -1665,7 +2187,7 @@ describe('Pi package executable-code boundary', () => { expect(state.disabledSources).toEqual([]); }); - it('keeps a snapshot timeout disabled across cache expiry until staging succeeds', async () => { + it('retries a transient snapshot duration limit on the next session', async () => { const { root, source } = await createPackage(); const now = Date.now(); const nowSpy = vi.spyOn(Date, 'now').mockReturnValue(now); @@ -1688,16 +2210,15 @@ describe('Pi package executable-code boundary', () => { promptTemplates: [], packageRoots: [], }); - expect(listener).toHaveBeenCalledTimes(1); + expect(listener).not.toHaveBeenCalled(); - // The failed staging started from a fresh inspection whose one-second - // cache is now stale. A Renderer refresh must still see the projected - // failure instead of rebuilding an enabled view from raw inspection. + // Wall-clock exhaustion depends on current machine load. It fails this + // session closed but must not become a durable package disable. nowSpy.mockReturnValue(now + 2_000); await expect(store.listPiPackages()).resolves.toMatchObject({ - packages: [{ source, enabled: false, warning: 'inspection-limit' }], + packages: [{ source, enabled: true }], }); - await expect(store.listManagedPiPromptCommands()).resolves.toEqual([]); + await expect(store.listManagedPiPromptCommands()).resolves.not.toEqual([]); const recoveredSnapshotRoot = path.join(runtime.userData, 'recovered-package-snapshot'); await expect(store.resolveManagedPiPackageResources({ @@ -1711,7 +2232,7 @@ describe('Pi package executable-code boundary', () => { extensions: [path.join(recoveredSnapshotRoot, '0', 'extensions', 'index.ts')], packageRoots: [path.join(recoveredSnapshotRoot, '0')], }); - expect(listener).toHaveBeenCalledTimes(2); + expect(listener).not.toHaveBeenCalled(); nowSpy.mockReturnValue(now + 4_000); const recoveredList = await store.listPiPackages(); @@ -1723,8 +2244,9 @@ describe('Pi package executable-code boundary', () => { const state = JSON.parse(await fs.readFile( path.join(runtime.userData, 'pi-package-home', 'cindy-package-state.json'), 'utf8', - )) as { disabledSources: string[] }; + )) as { disabledSources: string[]; snapshotUnavailablePackages: unknown[] }; expect(state.disabledSources).toEqual([]); + expect(state.snapshotUnavailablePackages).toEqual([]); await expect(fs.readFile( path.join(root, 'extensions', 'index.ts'), 'utf8', @@ -1737,7 +2259,7 @@ describe('Pi package executable-code boundary', () => { } }); - it('shares snapshot failures across instances and clears them after a successful staging', async () => { + it('shares deterministic snapshot failures and clears them through a production update', async () => { const { source } = await createPackage(); const firstStore = await import('../pi-package-store.js'); await mutateAuthorized(firstStore, { action: 'set-enabled', source, enabled: true }); @@ -1751,9 +2273,9 @@ describe('Pi package executable-code boundary', () => { await expect(firstStore.resolveManagedPiPackageResources({ snapshotRoot: path.join(runtime.userData, 'first-instance-failed-snapshot'), snapshotLimits: { - maxEntries: 100, + maxEntries: 1, maxBytes: 1024 * 1024, - maxDurationMs: 0, + maxDurationMs: 10_000, }, })).resolves.toEqual({ extensions: [], @@ -1770,6 +2292,7 @@ describe('Pi package executable-code boundary', () => { const unsubscribeFirst = firstStore.onPiPackagesChanged(firstListener); await new Promise((resolve) => setTimeout(resolve, 50)); try { + await mutateAuthorized(secondStore, { action: 'update', source }); const recoveredRoot = path.join(runtime.userData, 'second-instance-recovered-snapshot'); await expect(secondStore.resolveManagedPiPackageResources({ snapshotRoot: recoveredRoot, @@ -1784,8 +2307,8 @@ describe('Pi package executable-code boundary', () => { const state = JSON.parse(await fs.readFile( path.join(runtime.userData, 'pi-package-home', 'cindy-package-state.json'), 'utf8', - )) as { snapshotUnavailableRoots: Record }; - expect(state.snapshotUnavailableRoots).toEqual({}); + )) as { snapshotUnavailablePackages: unknown[] }; + expect(state.snapshotUnavailablePackages).toEqual([]); } finally { unsubscribeFirst(); } @@ -1904,16 +2427,16 @@ describe('Pi package executable-code boundary', () => { disabledSources: string[]; approvedExtensionSources: string[]; approvedExtensionFingerprints: Record; - snapshotUnavailableRoots: Record; + snapshotUnavailablePackages: unknown[]; }; expect(migrated).toEqual({ - version: 3, + version: 4, disabledSources: ['npm:keep-disabled'], approvedExtensionSources: [source], approvedExtensionFingerprints: { [source]: expect.stringMatching(/^[a-f0-9]{64}$/), }, - snapshotUnavailableRoots: {}, + snapshotUnavailablePackages: [], }); await mutateAuthorized(store, { action: 'install', source }); @@ -1923,6 +2446,38 @@ describe('Pi package executable-code boundary', () => { expect(installSpawn?.args).toContain('--no-approve'); }); + it('drops the v3 root-only snapshot failure while migrating durable package state', async () => { + const { root, source } = await createSkillOnlyPackage('npm:v3-snapshot-migration'); + const stateDir = path.join(runtime.userData, 'pi-package-home'); + await fs.mkdir(stateDir, { recursive: true }); + await fs.writeFile( + path.join(stateDir, 'cindy-package-state.json'), + JSON.stringify({ + version: 3, + disabledSources: [], + approvedExtensionSources: [], + approvedExtensionFingerprints: {}, + snapshotUnavailableRoots: { + [await fs.realpath(root)]: 'inspection-limit', + }, + }), + ); + + const store = await import('../pi-package-store.js'); + await store.mutatePiPackage({ action: 'set-enabled', source, enabled: false }); + + const migrated = JSON.parse( + await fs.readFile(path.join(stateDir, 'cindy-package-state.json'), 'utf8'), + ) as { + version: number; + disabledSources: string[]; + snapshotUnavailablePackages: unknown[]; + }; + expect(migrated.version).toBe(4); + expect(migrated.disabledSources).toEqual([source]); + expect(migrated.snapshotUnavailablePackages).toEqual([]); + }); + it.each([ ['transient I/O', 'EIO'], ['permission', 'EACCES'], @@ -2753,4 +3308,63 @@ describe('Pi package executable-code boundary', () => { }); expect(result.packages[129]?.warning).toBe('inspection-limit'); }); + + it('logs redacted correlated timings for every package startup stage', async () => { + const { root, source } = await createSkillOnlyPackage('npm:startup-timing'); + const snapshotRoot = path.join(runtime.userData, 'startup-timing-snapshot'); + const startupTraceId = '0123456789abcdef'; + const store = await import('../pi-package-store.js'); + + await expect(store.resolveManagedPiPackageResources({ + snapshotRoot, + startupTraceId, + })).resolves.toMatchObject({ skills: [expect.any(Object)] }); + + const events = loggerRuntime.info.mock.calls + .filter(([message]) => message === 'pi startup stage') + .map(([, fields]) => fields as Record); + expect(events.map((event) => event.stage)).toEqual([ + 'package-list', + 'package-inspection', + 'package-compatibility', + 'package-fingerprint', + 'package-snapshot', + ]); + for (const event of events) { + expect(event).toEqual({ + startupTraceId, + stage: expect.any(String), + durationMs: expect.any(Number), + status: 'ok', + packageCount: 1, + resourceCount: 1, + skippedPackageCount: 0, + }); + } + const serialized = JSON.stringify(events); + expect(serialized).not.toContain(root); + expect(serialized).not.toContain(source); + expect(serialized).not.toContain(snapshotRoot); + }); + + it('marks startup timing degraded when snapshot limits quarantine a package', async () => { + await createSkillOnlyPackage('npm:startup-timing-limited'); + const store = await import('../pi-package-store.js'); + + await store.resolveManagedPiPackageResources({ + snapshotRoot: path.join(runtime.userData, 'startup-timing-limited-snapshot'), + snapshotLimits: { maxEntries: 1, maxBytes: 1024 * 1024, maxDurationMs: 10_000 }, + startupTraceId: 'fedcba9876543210', + }); + + const snapshotEvent = loggerRuntime.info.mock.calls + .filter(([message]) => message === 'pi startup stage') + .map(([, fields]) => fields as Record) + .find((fields) => fields.stage === 'package-snapshot'); + expect(snapshotEvent).toMatchObject({ + startupTraceId: 'fedcba9876543210', + status: 'degraded', + skippedPackageCount: 1, + }); + }); }); diff --git a/apps/desktop/src/main/maker-host/pi-package-store.ts b/apps/desktop/src/main/maker-host/pi-package-store.ts index 1438f6537a..d5ddbadf40 100644 --- a/apps/desktop/src/main/maker-host/pi-package-store.ts +++ b/apps/desktop/src/main/maker-host/pi-package-store.ts @@ -73,13 +73,14 @@ const MAX_INSPECTED_PACKAGES = 128; const MAX_ALL_INSPECTION_MS = 10_000; const MAX_EXTENSION_FILES = 128; const INSPECTION_CACHE_MS = 1_000; +const SNAPSHOT_UNAVAILABLE_RETRY_MS = 6 * 60 * 60 * 1_000; const SNAPSHOT_COPY_CHUNK_BYTES = 256 * 1024; const DEFAULT_SNAPSHOT_LIMITS: PiPackageSnapshotLimits = { maxEntries: 10_000, maxBytes: 128 * 1024 * 1024, maxDurationMs: 15_000, }; -const STATE_VERSION = 3; +const STATE_VERSION = 4; const CHANGE_TOKEN_POLL_MS = 250; const changeListeners = new Set<() => void>(); let changeTokenWatcherActive = false; @@ -176,13 +177,37 @@ function stopPiPackageChangeTokenWatcher(): void { } type SnapshotUnavailableWarning = 'inspection-failed' | 'inspection-limit'; +type PiPackageSnapshotLimitReason = 'entries' | 'bytes' | 'duration'; + +interface SnapshotUnavailablePackageIdentity { + source: string; + installedRoot: string; + installationIdentity: string; + snapshotRoot: string; +} + +interface SnapshotUnavailablePackageProjection extends SnapshotUnavailablePackageIdentity { + warning: SnapshotUnavailableWarning; + retryAfterEpochMs?: number; +} + +/** Durable deterministic failure bound to one exact package installation. */ +interface PersistedSnapshotUnavailablePackage extends SnapshotUnavailablePackageIdentity { + warning: 'inspection-limit'; + retryAfterEpochMs: number; +} + +interface SnapshotUnavailableRootProjection { + warning: SnapshotUnavailableWarning; + durable: boolean; +} interface PiPackageState { version: typeof STATE_VERSION; disabledSources: string[]; approvedExtensionSources: string[]; approvedExtensionFingerprints: Record; - snapshotUnavailableRoots: Record; + snapshotUnavailablePackages: PersistedSnapshotUnavailablePackage[]; } type PiPackageStateReadResult = @@ -234,6 +259,79 @@ export interface PiPackageSnapshotLimits { maxDurationMs: number; } +const PI_PACKAGE_STARTUP_STAGES = [ + 'package-list', + 'package-inspection', + 'package-compatibility', + 'package-fingerprint', + 'package-snapshot', +] as const; +type PiPackageStartupStage = typeof PI_PACKAGE_STARTUP_STAGES[number]; + +interface PiPackageStartupTiming { + startupTraceId: string; + durationsMs: Record; + packageCount: number; + resourceCount: number; + skippedPackageCount: number; + degraded: boolean; +} + +function createPiPackageStartupTiming(startupTraceId: string | undefined): PiPackageStartupTiming | undefined { + if (!startupTraceId || !/^[a-f0-9]{16}$/.test(startupTraceId)) return undefined; + return { + startupTraceId, + durationsMs: { + 'package-list': 0, + 'package-inspection': 0, + 'package-compatibility': 0, + 'package-fingerprint': 0, + 'package-snapshot': 0, + }, + packageCount: 0, + resourceCount: 0, + skippedPackageCount: 0, + degraded: false, + }; +} + +function recordPiPackageStartupDuration( + timing: PiPackageStartupTiming | undefined, + stage: PiPackageStartupStage, + startedAt: number, +): void { + if (!timing) return; + timing.durationsMs[stage] += Math.max(0, Date.now() - startedAt); +} + +async function measurePiPackageStartupStage( + timing: PiPackageStartupTiming | undefined, + stage: PiPackageStartupStage, + operation: () => Promise, +): Promise { + const startedAt = Date.now(); + try { + return await operation(); + } finally { + recordPiPackageStartupDuration(timing, stage, startedAt); + } +} + +function emitPiPackageStartupTiming(timing: PiPackageStartupTiming | undefined): void { + if (!timing) return; + for (const stage of PI_PACKAGE_STARTUP_STAGES) { + log.info('pi startup stage', { + startupTraceId: timing.startupTraceId, + stage, + durationMs: timing.durationsMs[stage], + status: timing.degraded ? 'degraded' : 'ok', + packageCount: timing.packageCount, + resourceCount: timing.resourceCount, + skippedPackageCount: timing.skippedPackageCount, + }); + } +} + interface InspectedPackage { /** Original Pi-owned identifier. Never expose this field across IPC. */ rawSource: string; @@ -246,6 +344,10 @@ interface InspectedPackage { contentFingerprint?: string; /** Persisted approval exists but no longer matches the current package tree. */ staleApproval?: boolean; + /** Deterministic snapshot failure reused without walking the package tree. */ + snapshotUnavailable?: PersistedSnapshotUnavailablePackage; + installationIdentity?: string; + snapshotRoot?: string; } interface PackageSourceProjection { @@ -267,7 +369,7 @@ interface InspectionBudget { } class PiPackageInspectionLimitError extends Error { - constructor() { + constructor(readonly reason: PiPackageSnapshotLimitReason = 'entries') { super('Pi package inspection limit exceeded'); this.name = 'PiPackageInspectionLimitError'; } @@ -277,7 +379,7 @@ let mutationTail: Promise = Promise.resolve(); let inspectionPromise: Promise | undefined; let inspectionCache: { expiresAt: number; value: InspectedPackage[] } | undefined; let inspectionGeneration = 0; -const snapshotUnavailableRoots = new Map(); +const snapshotUnavailablePackages = new Map(); function packageHome(): string { return path.join(app.getPath('userData'), 'pi-package-home'); @@ -304,6 +406,44 @@ async function snapshotRootForInstalledPackage( } } +async function packageInstallationIdentity( + installedRoot: string, + stat: Stats, +): Promise { + let manifestIdentity: readonly string[] = ['not-directory']; + if (stat.isDirectory()) { + try { + const manifest = await readUtf8FileBounded( + path.join(installedRoot, 'package.json'), + MAX_PACKAGE_JSON_BYTES, + installedRoot, + ); + manifestIdentity = [ + 'package-json', + createHash('sha256').update(manifest.text).digest('hex'), + ]; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; + manifestIdentity = ['package-json-missing']; + } + } + return createHash('sha256') + .update( + JSON.stringify([ + 'cindy-pi-package-installation-v2', + stat.dev, + stat.ino, + stat.mode, + stat.size, + stat.birthtimeMs, + stat.mtimeMs, + stat.ctimeMs, + manifestIdentity, + ]), + ) + .digest('hex'); +} + function statePath(): string { return path.join(packageHome(), 'cindy-package-state.json'); } @@ -356,31 +496,78 @@ function parseApprovedExtensionFingerprints(value: unknown): Record; } -function parseSnapshotUnavailableRoots( +function parseSnapshotUnavailablePackages( value: unknown, -): Record | undefined { - if (value === undefined) return {}; - if (!value || typeof value !== 'object' || Array.isArray(value)) return undefined; - const entries = Object.entries(value); +): PersistedSnapshotUnavailablePackage[] | undefined { + if (value === undefined) return []; + if (!Array.isArray(value)) return undefined; + const entries = value as Array>; if ( entries.length > MAX_INSPECTED_PACKAGES - || !entries.every(([root, warning]) => ( - root.length > 0 - && root.length <= MAX_SOURCE_LENGTH - && (warning === 'inspection-failed' || warning === 'inspection-limit') + || !entries.every((entry) => ( + typeof entry.source === 'string' + && entry.source.length > 0 + && entry.source.length <= MAX_SOURCE_LENGTH + && typeof entry.installedRoot === 'string' + && entry.installedRoot.length > 0 + && entry.installedRoot.length <= MAX_SOURCE_LENGTH + && typeof entry.snapshotRoot === 'string' + && entry.snapshotRoot.length > 0 + && entry.snapshotRoot.length <= MAX_SOURCE_LENGTH + && typeof entry.installationIdentity === 'string' + && /^[a-f0-9]{64}$/.test(entry.installationIdentity) + && entry.warning === 'inspection-limit' + && ( + entry.retryAfterEpochMs === undefined + || ( + typeof entry.retryAfterEpochMs === 'number' + && Number.isSafeInteger(entry.retryAfterEpochMs) + && entry.retryAfterEpochMs >= 0 + ) + ) )) ) return undefined; - return Object.fromEntries( - entries.sort(([left], [right]) => left.localeCompare(right)), - ) as Record; + return sortSnapshotUnavailablePackages(entries.map((entry) => ({ + ...entry, + // v4 states written before this additive field retry once after upgrade. + retryAfterEpochMs: entry.retryAfterEpochMs ?? 0, + })) as PersistedSnapshotUnavailablePackage[]); +} + +function sortSnapshotUnavailablePackages( + entries: PersistedSnapshotUnavailablePackage[], +): PersistedSnapshotUnavailablePackage[] { + return entries.toSorted((left, right) => ( + left.source.localeCompare(right.source) + || left.installedRoot.localeCompare(right.installedRoot) + || left.snapshotRoot.localeCompare(right.snapshotRoot) + )); } -function applySharedSnapshotUnavailableRoots( - roots: Readonly>, +function isSnapshotUnavailableRetryDue( + entry: SnapshotUnavailablePackageProjection, +): boolean { + return entry.retryAfterEpochMs !== undefined && Date.now() >= entry.retryAfterEpochMs; +} + +function snapshotUnavailablePackageKey( + entry: SnapshotUnavailablePackageIdentity, +): string { + return JSON.stringify([ + entry.source, + path.resolve(entry.installedRoot), + entry.installationIdentity, + path.resolve(entry.snapshotRoot), + ]); +} + +function applySharedSnapshotUnavailablePackages( + packages: readonly SnapshotUnavailablePackageProjection[], ): void { - snapshotUnavailableRoots.clear(); - for (const [root, warning] of Object.entries(roots)) { - snapshotUnavailableRoots.set(path.resolve(root), warning); + snapshotUnavailablePackages.clear(); + for (const pkg of packages) { + if (isSnapshotUnavailableRetryDue(pkg)) continue; + snapshotUnavailablePackages.set(snapshotUnavailablePackageKey(pkg), pkg); } } @@ -390,7 +577,7 @@ function emptyState(): PiPackageState { disabledSources: [], approvedExtensionSources: [], approvedExtensionFingerprints: {}, - snapshotUnavailableRoots: {}, + snapshotUnavailablePackages: [], }; } @@ -400,7 +587,7 @@ async function readState(): Promise { const fingerprints = parseApprovedExtensionFingerprints( parsed.approvedExtensionFingerprints, ); - const unavailableRoots = parseSnapshotUnavailableRoots(parsed.snapshotUnavailableRoots); + const unavailablePackages = parseSnapshotUnavailablePackages(parsed.snapshotUnavailablePackages); if ( parsed.version === STATE_VERSION && Array.isArray(parsed.disabledSources) @@ -408,7 +595,7 @@ async function readState(): Promise { && Array.isArray(parsed.approvedExtensionSources) && parsed.approvedExtensionSources.every((source) => typeof source === 'string') && fingerprints - && unavailableRoots + && unavailablePackages ) { const approvedExtensionSources = [...new Set(parsed.approvedExtensionSources)] .filter((source) => Object.hasOwn(fingerprints, source)); @@ -421,25 +608,41 @@ async function readState(): Promise { approvedExtensionFingerprints: Object.fromEntries( approvedExtensionSources.map((source) => [source, fingerprints[source]!]), ), - snapshotUnavailableRoots: unavailableRoots, + snapshotUnavailablePackages: unavailablePackages, }, }; } if ( - (parsed.version === 1 || parsed.version === 2) + (parsed.version === 1 || parsed.version === 2 || parsed.version === 3) && Array.isArray(parsed.disabledSources) && parsed.disabledSources.every((source) => typeof source === 'string') ) { - // Preserve explicit disables. Older approvals had no byte identity, so - // they cannot authorize executable code under the v3 content boundary. + const legacyFingerprints = + parsed.version === 3 + ? parseApprovedExtensionFingerprints(parsed.approvedExtensionFingerprints) + : undefined; + const legacyApprovedSources = + parsed.version === 3 && + Array.isArray(parsed.approvedExtensionSources) && + parsed.approvedExtensionSources.every((source) => typeof source === 'string') && + legacyFingerprints + ? [...new Set(parsed.approvedExtensionSources)].filter((source) => + Object.hasOwn(legacyFingerprints, source), + ) + : []; + // v3 approvals already carry byte identity and remain valid. Its root-only + // negative cache cannot prove source/install identity, so migration drops it. + // v1/v2 approvals had no byte identity and stay fail closed as before. return { ok: true, state: { version: STATE_VERSION, disabledSources: [...new Set(parsed.disabledSources)], - approvedExtensionSources: [], - approvedExtensionFingerprints: {}, - snapshotUnavailableRoots: {}, + approvedExtensionSources: legacyApprovedSources, + approvedExtensionFingerprints: Object.fromEntries( + legacyApprovedSources.map((source) => [source, legacyFingerprints![source]!]), + ), + snapshotUnavailablePackages: [], }, }; } @@ -669,12 +872,11 @@ function createInspectionBudget(): InspectionBudget { function assertInspectionBudget(budget: InspectionBudget, depth = 0, increment = 0): void { budget.entries += increment; - if ( - depth > MAX_INSPECTION_DEPTH - || budget.entries > MAX_INSPECTION_ENTRIES - || Date.now() - budget.startedAt > MAX_INSPECTION_MS - ) { - throw new PiPackageInspectionLimitError(); + if (depth > MAX_INSPECTION_DEPTH || budget.entries > MAX_INSPECTION_ENTRIES) { + throw new PiPackageInspectionLimitError('entries'); + } + if (Date.now() - budget.startedAt > MAX_INSPECTION_MS) { + throw new PiPackageInspectionLimitError('duration'); } } @@ -690,7 +892,7 @@ async function readUtf8FileBounded( 'Pi package metadata changed before reading', ); try { - if (stat.size > maxBytes) throw new PiPackageInspectionLimitError(); + if (stat.size > maxBytes) throw new PiPackageInspectionLimitError('bytes'); const buffer = Buffer.alloc(maxBytes + 1); let bytes = 0; while (bytes < buffer.length) { @@ -698,7 +900,7 @@ async function readUtf8FileBounded( if (result.bytesRead === 0) break; bytes += result.bytesRead; } - if (bytes > maxBytes) throw new PiPackageInspectionLimitError(); + if (bytes > maxBytes) throw new PiPackageInspectionLimitError('bytes'); const after = await handle.stat(); if (!sameStableFileIdentity(stat, after) || bytes !== after.size) { throw new Error('Pi package metadata changed while reading'); @@ -715,7 +917,7 @@ async function readInspectionMetadata( confinementRoot: string, ): Promise { const remaining = MAX_INSPECTION_METADATA_BYTES - budget.metadataBytes; - if (remaining < 0) throw new PiPackageInspectionLimitError(); + if (remaining < 0) throw new PiPackageInspectionLimitError('bytes'); const result = await readUtf8FileBounded(file, remaining, confinementRoot); budget.metadataBytes += result.bytes; assertInspectionBudget(budget); @@ -725,7 +927,7 @@ async function readInspectionMetadata( function normalizeManifestEntries(value: unknown, fallback: string[]): string[] { if (value === undefined) return fallback; if (!Array.isArray(value)) throw new Error('Invalid Pi package manifest entries'); - if (value.length > MAX_MANIFEST_ENTRIES) throw new PiPackageInspectionLimitError(); + if (value.length > MAX_MANIFEST_ENTRIES) throw new PiPackageInspectionLimitError('entries'); const entries: string[] = []; for (const entry of value) { if ( @@ -1069,7 +1271,12 @@ function resourceView(kind: Exclude, file: s }; } -async function extensionResourceView(root: string, file: string): Promise { +async function extensionResourceView( + root: string, + file: string, + startupTiming?: PiPackageStartupTiming, +): Promise { + const startedAt = Date.now(); try { const analysis = await analyzePiExtensionCompatibility(file, root); return { @@ -1088,6 +1295,8 @@ async function extensionResourceView(root: string, file: string): Promise>, aggregateBudget: SnapshotBudgetCounters, + startupTiming?: PiPackageStartupTiming, ): Promise { const current = cache.get(root); if (current) return current; - const pending = fingerprintPiPackageTree(root, DEFAULT_SNAPSHOT_LIMITS, aggregateBudget); + const pending = measurePiPackageStartupStage( + startupTiming, + 'package-fingerprint', + () => fingerprintPiPackageTree(root, DEFAULT_SNAPSHOT_LIMITS, aggregateBudget), + ); cache.set(root, pending); return pending; } @@ -1158,6 +1372,7 @@ async function inspectPackage( state: PiPackageState, fingerprintCache: Map>, aggregateFingerprintBudget: SnapshotBudgetCounters, + startupTiming?: PiPackageStartupTiming, ): Promise { const empty: PiManagedPackageResources = { extensions: [], skills: [], promptTemplates: [], packageRoots: [], @@ -1195,6 +1410,8 @@ async function inspectPackage( }; } let installedRoot: string | undefined; + let installationIdentity: string | undefined; + let snapshotRoot: string | undefined; try { const budget = createInspectionBudget(); const { canonicalPath: root, stat: rootStat } = await resolveStablePackagePath( @@ -1202,6 +1419,8 @@ async function inspectPackage( 'Pi package root changed during inspection', ); installedRoot = root; + installationIdentity = await packageInstallationIdentity(root, rootStat); + snapshotRoot = await snapshotRootForInstalledPackage(pkg.source, root); if (pkg.filtered) { return { rawSource: pkg.source, @@ -1216,17 +1435,21 @@ async function inspectPackage( launch: empty, promptCommands: [], installedRoot: root, + installationIdentity, + snapshotRoot, }; } if (rootStat.isFile()) { const isExtension = /\.(?:ts|js)$/i.test(root); - const launchRoot = await snapshotRootForInstalledPackage(pkg.source, root); - const resources = isExtension ? [await extensionResourceView(path.dirname(root), root)] : []; + const resources = isExtension + ? [await extensionResourceView(path.dirname(root), root, startupTiming)] + : []; const contentFingerprint = isExtension ? await fingerprintPackageTreeCached( - launchRoot, + snapshotRoot, fingerprintCache, aggregateFingerprintBudget, + startupTiming, ) : undefined; const requiresExtensionApproval = isExtension && !( @@ -1249,10 +1472,12 @@ async function inspectPackage( ...(resources.length === 0 ? { warning: 'no-resources' as const } : {}), }, launch: enabled && isExtension - ? { extensions: [root], skills: [], promptTemplates: [], packageRoots: [launchRoot] } + ? { extensions: [root], skills: [], promptTemplates: [], packageRoots: [snapshotRoot] } : empty, promptCommands: [], installedRoot: root, + installationIdentity, + snapshotRoot, ...(contentFingerprint ? { contentFingerprint } : {}), ...(staleApproval ? { staleApproval: true } : {}), }; @@ -1266,16 +1491,22 @@ async function inspectPackage( } catch (error) { if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; } - const runtimeRequirements = evaluatePiRuntimeRequirements( - manifest.peerDependencies, - await getCurrentPiVersion(), - ).map((requirement) => ({ - ...requirement, - range: truncateDisplayField(requirement.range, MAX_DISPLAY_NAME_BYTES), - ...(requirement.currentVersion - ? { currentVersion: truncateDisplayField(requirement.currentVersion, MAX_DISPLAY_VERSION_BYTES) } - : {}), - })); + const compatibilityStartedAt = Date.now(); + let runtimeRequirements; + try { + runtimeRequirements = evaluatePiRuntimeRequirements( + manifest.peerDependencies, + await getCurrentPiVersion(), + ).map((requirement) => ({ + ...requirement, + range: truncateDisplayField(requirement.range, MAX_DISPLAY_NAME_BYTES), + ...(requirement.currentVersion + ? { currentVersion: truncateDisplayField(requirement.currentVersion, MAX_DISPLAY_VERSION_BYTES) } + : {}), + })); + } finally { + recordPiPackageStartupDuration(startupTiming, 'package-compatibility', compatibilityStartedAt); + } const declared = manifest.pi; const extensionEntries = normalizeManifestEntries(declared?.extensions, ['extensions']); const skillEntries = normalizeManifestEntries(declared?.skills, ['skills']); @@ -1292,14 +1523,16 @@ async function inspectPackage( collectFilesByExtension(await confinedExistingPaths(root, themeInputs), ['.json'], budget), ]); assertInspectionBudget(budget); - if (extensions.length > MAX_EXTENSION_FILES) throw new PiPackageInspectionLimitError(); + if (extensions.length > MAX_EXTENSION_FILES) { + throw new PiPackageInspectionLimitError('entries'); + } // Babel parsing happens in Electron's main process. Keep analysis // sequential and re-check the package-wide wall-clock budget between // entries so a package cannot fan out thousands of CPU-heavy parses. const extensionResources: PiPackageResourceView[] = []; for (const file of extensions) { assertInspectionBudget(budget); - extensionResources.push(await extensionResourceView(root, file)); + extensionResources.push(await extensionResourceView(root, file, startupTiming)); assertInspectionBudget(budget); } const resources: PiPackageResourceView[] = [ @@ -1308,7 +1541,6 @@ async function inspectPackage( ...prompts.map((file) => resourceView('prompt', file)), ...themes.map((file) => resourceView('theme', file)), ]; - const launchRoot = await snapshotRootForInstalledPackage(pkg.source, root); const hasLaunchResources = extensions.length > 0 || skills.length > 0 || prompts.length > 0; // Every enabled directory package is copied as one launch root, including // Skills/Prompts-only packages. Apply the exact snapshot tree limits here @@ -1316,9 +1548,10 @@ async function inspectPackage( // aborting the combined task snapshot and hiding otherwise valid packages. const contentFingerprint = hasLaunchResources ? await fingerprintPackageTreeCached( - launchRoot, + snapshotRoot, fingerprintCache, aggregateFingerprintBudget, + startupTiming, ) : undefined; const requiresExtensionApproval = extensions.length > 0 && !( @@ -1357,10 +1590,12 @@ async function inspectPackage( ...(warning ? { warning } : {}), }, launch: enabled && hasLaunchResources - ? { extensions, skills, promptTemplates: prompts, packageRoots: [launchRoot] } + ? { extensions, skills, promptTemplates: prompts, packageRoots: [snapshotRoot] } : empty, promptCommands, installedRoot: root, + installationIdentity, + snapshotRoot, ...(contentFingerprint ? { contentFingerprint } : {}), ...(staleApproval ? { staleApproval: true } : {}), }; @@ -1369,6 +1604,20 @@ async function inspectPackage( source: displaySource, message: error instanceof Error ? error.message : String(error), }); + const deterministicLimit = isDurableInspectionFailure(error); + const snapshotUnavailable = deterministicLimit + && installedRoot + && installationIdentity + && snapshotRoot + ? { + source: pkg.source, + installedRoot, + installationIdentity, + snapshotRoot, + warning: 'inspection-limit' as const, + retryAfterEpochMs: Date.now() + SNAPSHOT_UNAVAILABLE_RETRY_MS, + } + : undefined; return { rawSource: pkg.source, view: { @@ -1385,29 +1634,86 @@ async function inspectPackage( launch: empty, promptCommands: [], ...(installedRoot ? { installedRoot } : {}), + ...(installationIdentity ? { installationIdentity } : {}), + ...(snapshotRoot ? { snapshotRoot } : {}), + ...(snapshotUnavailable ? { snapshotUnavailable } : {}), + }; + } +} + +async function persistedSnapshotLimitProjection( + pkg: ListedPackage, + state: PiPackageState, +): Promise { + if (!pkg.installedPath || pkg.filtered) return undefined; + if (state.snapshotUnavailablePackages.length === 0) return undefined; + const { displaySource, unsafe } = projectPackageSource(pkg.source); + if (unsafe) return undefined; + try { + const { canonicalPath: installedRoot, stat } = await resolveStablePackagePath( + pkg.installedPath, + 'Pi package root changed while resolving its snapshot failure', + ); + const snapshotRoot = await snapshotRootForInstalledPackage(pkg.source, installedRoot); + const installationIdentity = await packageInstallationIdentity(installedRoot, stat); + const persisted = state.snapshotUnavailablePackages.find( + (entry) => + entry.source === pkg.source && + path.resolve(entry.installedRoot) === path.resolve(installedRoot) && + path.resolve(entry.snapshotRoot) === path.resolve(snapshotRoot) && + entry.installationIdentity === installationIdentity, + ); + if (!persisted || isSnapshotUnavailableRetryDue(persisted)) return undefined; + return { + rawSource: pkg.source, + view: { + source: displaySource, + name: displaySource, + enabled: false, + canToggle: false, + resources: [], + warning: 'inspection-limit', + }, + launch: { extensions: [], skills: [], promptTemplates: [], packageRoots: [] }, + promptCommands: [], + installedRoot, + installationIdentity, + snapshotRoot, + snapshotUnavailable: persisted, }; + } catch { + // Path identity could not be proven cheaply. Fall through to the regular + // fail-closed inspection so transient filesystem failures remain retryable. + return undefined; } } -async function inspectAllPackagesUncached(): Promise { +async function inspectAllPackagesUncached( + options: { + ignorePersistedSnapshotUnavailable?: boolean; + startupTiming?: PiPackageStartupTiming; + } = {}, +): Promise { const [{ stdout }, stateResult] = await Promise.all([ - runPiPackageCommand( - ['list', '--no-approve'], - COMMAND_TIMEOUT_MS, - { requireCompleteStdout: true }, + measurePiPackageStartupStage( + options.startupTiming, + 'package-list', + () => runPiPackageCommand( + ['list', '--no-approve'], + COMMAND_TIMEOUT_MS, + { requireCompleteStdout: true }, + ), ), readState(), ]); - // Snapshot failures are shared package-store state, not a property of one - // Main process. Every fresh inspection replaces the local projection with - // the atomically persisted view so packaged/dev peers agree after the - // existing change-token invalidation. const state = stateResult.ok ? stateResult.state : emptyState(); if (stateResult.ok) { - applySharedSnapshotUnavailableRoots(state.snapshotUnavailableRoots); + applySharedSnapshotUnavailablePackages(state.snapshotUnavailablePackages); } const listed = parsePiPackageListOutput(stdout); + if (options.startupTiming) options.startupTiming.packageCount = listed.length; const startedAt = Date.now(); + const inspectionStartedAt = Date.now(); const inspected: InspectedPackage[] = []; const fingerprintCache = new Map>(); const aggregateFingerprintBudget = createSnapshotBudgetCounters(DEFAULT_SNAPSHOT_LIMITS); @@ -1429,12 +1735,23 @@ async function inspectAllPackagesUncached(): Promise { }); continue; } - const inspectedPackage = await inspectPackage( - pkg, - state, - fingerprintCache, - aggregateFingerprintBudget, - ); + const inspectedPackage = + stateResult.ok && options.ignorePersistedSnapshotUnavailable !== true + ? ((await persistedSnapshotLimitProjection(pkg, state)) ?? + (await inspectPackage( + pkg, + state, + fingerprintCache, + aggregateFingerprintBudget, + options.startupTiming, + ))) + : await inspectPackage( + pkg, + state, + fingerprintCache, + aggregateFingerprintBudget, + options.startupTiming, + ); if (stateResult.ok) { inspected.push(inspectedPackage); } else { @@ -1454,12 +1771,19 @@ async function inspectAllPackagesUncached(): Promise { ...(inspectedPackage.contentFingerprint ? { contentFingerprint: inspectedPackage.contentFingerprint } : {}), + ...(inspectedPackage.installationIdentity + ? { installationIdentity: inspectedPackage.installationIdentity } + : {}), + ...(inspectedPackage.snapshotRoot + ? { snapshotRoot: inspectedPackage.snapshotRoot } + : {}), }); } // Package inspection includes synchronous parser work in Electron's main // process. Yield between packages so a long roster cannot monopolize it. await new Promise((resolve) => setImmediate(resolve)); } + recordPiPackageStartupDuration(options.startupTiming, 'package-inspection', inspectionStartedAt); return inspected; } @@ -1485,7 +1809,12 @@ async function inspectAllPackages(): Promise { return pending; } -async function inspectAllPackagesFreshUnderMutationLock(): Promise { +async function inspectAllPackagesFreshUnderMutationLock( + options: { + ignorePersistedSnapshotUnavailable?: boolean; + startupTiming?: PiPackageStartupTiming; + } = {}, +): Promise { // A local inspection that began before another process changed the shared // package store must finish before its generation is retired. Starting the // replacement under the cross-process mutation lock then re-reads both the @@ -1493,7 +1822,9 @@ async function inspectAllPackagesFreshUnderMutationLock(): Promise undefined); invalidateInspectionCache(); - return inspectAllPackages(); + return options.ignorePersistedSnapshotUnavailable || options.startupTiming + ? inspectAllPackagesUncached(options) + : inspectAllPackages(); } async function listPiPackagesNow(): Promise { @@ -1582,52 +1913,99 @@ export async function capturePiPackageEnableIdentity(source: string): Promise, + unavailableRoots: Iterable, + inspected: readonly InspectedPackage[] = [], ): Promise { const state = await requireState(); - const next: Record = {}; - for (const [root, warning] of unavailableRoots) { - next[path.resolve(root)] = warning; - } - const entries = Object.entries(next).sort(([left], [right]) => left.localeCompare(right)); - const currentEntries = Object.entries(state.snapshotUnavailableRoots) - .sort(([left], [right]) => left.localeCompare(right)); - const changed = entries.length !== currentEntries.length - || entries.some(([root, warning], index) => ( - root !== currentEntries[index]?.[0] || warning !== currentEntries[index]?.[1] - )); + const projectionsByRoot = new Map(); + for (const [root, projection] of unavailableRoots) { + projectionsByRoot.set(path.resolve(root), projection); + } + const projectionByIdentity = new Map(); + const durableByIdentity = new Map(); + for (const pkg of inspected) { + const retained = pkg.snapshotUnavailable; + if (retained) { + // A package-level inspection failure is already bound to its exact + // installation. Do not lower it to the shared npm resolver root, which + // would quarantine healthy sibling packages that happen to copy together. + const key = snapshotUnavailablePackageKey(retained); + projectionByIdentity.set(key, retained); + durableByIdentity.set(key, retained); + } + if (!pkg.installedRoot || !pkg.installationIdentity) continue; + for (const snapshotRoot of pkg.launch.packageRoots) { + const rootProjection = projectionsByRoot.get(path.resolve(snapshotRoot)); + if (!rootProjection) continue; + try { + const { canonicalPath, stat } = await resolveStablePackagePath( + pkg.installedRoot, + 'Pi package root changed while persisting its snapshot failure', + ); + if (path.resolve(canonicalPath) !== path.resolve(pkg.installedRoot)) continue; + const currentInstallationIdentity = await packageInstallationIdentity(canonicalPath, stat); + if (currentInstallationIdentity !== pkg.installationIdentity) continue; + const projection: SnapshotUnavailablePackageProjection = { + source: pkg.rawSource, + installedRoot: canonicalPath, + installationIdentity: currentInstallationIdentity, + snapshotRoot: path.resolve(snapshotRoot), + warning: rootProjection.warning, + }; + const key = snapshotUnavailablePackageKey(projection); + projectionByIdentity.set(key, projection); + if (rootProjection.warning === 'inspection-limit' && rootProjection.durable) { + durableByIdentity.set(key, { + ...projection, + warning: 'inspection-limit', + retryAfterEpochMs: Date.now() + SNAPSHOT_UNAVAILABLE_RETRY_MS, + }); + } + } catch { + // Identity cannot be proven, so the failure remains retryable. + } + } + } + const nextPackages = sortSnapshotUnavailablePackages([...durableByIdentity.values()]); + const changed = JSON.stringify(nextPackages) !== JSON.stringify(state.snapshotUnavailablePackages); if (changed) { await writeState({ ...state, - snapshotUnavailableRoots: Object.fromEntries(entries), + snapshotUnavailablePackages: nextPackages, }); } - applySharedSnapshotUnavailableRoots(next); + applySharedSnapshotUnavailablePackages([...projectionByIdentity.values()]); return changed; } function snapshotUnavailableWarningForPackage( pkg: InspectedPackage, ): SnapshotUnavailableWarning | undefined { - let warning: SnapshotUnavailableWarning | undefined; - for (const root of pkg.launch.packageRoots) { - const candidate = snapshotUnavailableRoots.get(path.resolve(root)); - if (candidate === 'inspection-failed') return candidate; - if (candidate) warning = candidate; - } - return warning; + if (pkg.snapshotUnavailable) return pkg.snapshotUnavailable.warning; + if (!pkg.installedRoot || !pkg.installationIdentity || !pkg.snapshotRoot) return undefined; + return snapshotUnavailablePackages.get(snapshotUnavailablePackageKey({ + source: pkg.rawSource, + installedRoot: pkg.installedRoot, + installationIdentity: pkg.installationIdentity, + snapshotRoot: pkg.snapshotRoot, + }))?.warning; } export async function resolveManagedPiPackageResources( - options?: { snapshotRoot: string; snapshotLimits?: PiPackageSnapshotLimits }, + options?: { + snapshotRoot: string; + snapshotLimits?: PiPackageSnapshotLimits; + startupTraceId?: string; + }, ): Promise { + const startupTiming = createPiPackageStartupTiming(options?.startupTraceId); if (!getReadyBinaryPath('pi')) { return { extensions: [], skills: [], promptTemplates: [], packageRoots: [] }; } try { const resolveResources = async (forceFresh = false): Promise => { const inspected = forceFresh - ? await inspectAllPackagesFreshUnderMutationLock() + ? await inspectAllPackagesFreshUnderMutationLock({ startupTiming }) : await inspectAllPackages(); if (options) { const staleApprovals = inspected @@ -1644,6 +2022,17 @@ export async function resolveManagedPiPackageResources( promptTemplates: [...new Set(inspected.flatMap((pkg) => pkg.launch.promptTemplates))], packageRoots: [...new Set(inspected.flatMap((pkg) => pkg.launch.packageRoots))], }; + if (startupTiming) { + startupTiming.resourceCount = resources.extensions.length + + resources.skills.length + + resources.promptTemplates.length; + startupTiming.skippedPackageCount = inspected.filter((pkg) => ( + pkg.view.warning === 'inspection-failed' + || pkg.view.warning === 'inspection-limit' + || Boolean(pkg.snapshotUnavailable) + )).length; + startupTiming.degraded = startupTiming.skippedPackageCount > 0; + } if (!options) return resources; const approvalsByRoot = new Map>(); @@ -1655,6 +2044,7 @@ export async function resolveManagedPiPackageResources( approvalsByRoot.set(root, approvals); } } + const snapshotStartedAt = Date.now(); try { const snapshotLimits = options.snapshotLimits ?? DEFAULT_SNAPSHOT_LIMITS; let staged = await stageManagedPackageSnapshot( @@ -1667,8 +2057,11 @@ export async function resolveManagedPiPackageResources( const copiedSourceRoots = stageMetadata?.sourcePackageRoots ?? resources.packageRoots; const verificationBudget = createSnapshotBudgetCounters(snapshotLimits); const unavailableVerificationRoots = new Map(); + const transientUnavailableVerificationRoots = new Set(); + const durableUnavailableVerificationRoots = new Set(); const failedVerificationIndexes = new Set(); let aggregateVerificationLimitReached = false; + let aggregateVerificationLimitReason: PiPackageSnapshotLimitReason | undefined; // Fingerprint verification has the same partial-success contract as // staging: a budget breach quarantines only the unverified roots, so // already copied and authenticated resources remain usable. @@ -1677,6 +2070,9 @@ export async function resolveManagedPiPackageResources( if (!approvals?.length) continue; if (aggregateVerificationLimitReached) { unavailableVerificationRoots.set(sourceRoot, 'inspection-limit'); + if (aggregateVerificationLimitReason === 'duration') { + transientUnavailableVerificationRoots.add(sourceRoot); + } failedVerificationIndexes.add(index); continue; } @@ -1692,8 +2088,17 @@ export async function resolveManagedPiPackageResources( } catch (error) { if (!(error instanceof PiPackageSnapshotLimitError)) throw error; unavailableVerificationRoots.set(sourceRoot, 'inspection-limit'); + if (error.reason === 'duration') { + transientUnavailableVerificationRoots.add(sourceRoot); + } + if (error.scope === 'package' && error.reason !== 'duration') { + durableUnavailableVerificationRoots.add(sourceRoot); + } failedVerificationIndexes.add(index); - if (error.scope === 'aggregate') aggregateVerificationLimitReached = true; + if (error.scope === 'aggregate') { + aggregateVerificationLimitReached = true; + aggregateVerificationLimitReason = error.reason; + } continue; } for (const approval of approvals) { @@ -1701,6 +2106,13 @@ export async function resolveManagedPiPackageResources( } } if (changedSources.size > 0) { + if (startupTiming) { + startupTiming.degraded = true; + startupTiming.skippedPackageCount = Math.max( + startupTiming.skippedPackageCount, + changedSources.size, + ); + } await fs.rm(options.snapshotRoot, { recursive: true, force: true }); await revokeExtensionApproval(changedSources); await publishPiPackagesChanged(); @@ -1733,24 +2145,56 @@ export async function resolveManagedPiPackageResources( ...(stageMetadata?.skippedPackageRoots ?? []), ...failedSources, ]; + const nextTransientSkippedRoots = [ + ...(stageMetadata?.transientSkippedPackageRoots ?? []), + ...transientUnavailableVerificationRoots, + ]; + const nextDurableSkippedRoots = [ + ...(stageMetadata?.durableSkippedPackageRoots ?? []), + ...durableUnavailableVerificationRoots, + ]; snapshotStageMetadata.set(filtered, { sourcePackageRoots: copiedSourceRoots.filter((_, index) => ( !failedVerificationIndexes.has(index) )), skippedPackageRoots: [...new Set(nextSkippedRoots)], + transientSkippedPackageRoots: [...new Set(nextTransientSkippedRoots)], + durableSkippedPackageRoots: [...new Set(nextDurableSkippedRoots)], }); staged = filtered; } - const unavailableRoots = new Map(); + const unavailableRoots = new Map(); + const transientSkippedRoots = new Set(stageMetadata?.transientSkippedPackageRoots ?? []); + const durableSkippedRoots = new Set(stageMetadata?.durableSkippedPackageRoots ?? []); for (const root of stageMetadata?.skippedPackageRoots ?? []) { - unavailableRoots.set(root, 'inspection-limit'); + if (transientSkippedRoots.has(root)) continue; + unavailableRoots.set(root, { + warning: 'inspection-limit', + durable: durableSkippedRoots.has(root), + }); } for (const [root, warning] of unavailableVerificationRoots) { - unavailableRoots.set(root, warning); + if (transientUnavailableVerificationRoots.has(root)) continue; + unavailableRoots.set(root, { + warning, + durable: durableUnavailableVerificationRoots.has(root), + }); + } + if (startupTiming) { + const skippedThisStart = new Set([ + ...(stageMetadata?.skippedPackageRoots ?? []), + ...unavailableVerificationRoots.keys(), + ]); + startupTiming.skippedPackageCount = Math.max( + startupTiming.skippedPackageCount, + skippedThisStart.size, + ); + startupTiming.degraded ||= skippedThisStart.size > 0; } const snapshotProjectionChanged = await persistSnapshotUnavailableProjection( unavailableRoots, + inspected, ); if (snapshotProjectionChanged) { await publishPiPackagesChanged({ invalidateCache: false }); @@ -1761,22 +2205,33 @@ export async function resolveManagedPiPackageResources( const warning = error instanceof PiPackageSnapshotLimitError ? 'inspection-limit' : 'inspection-failed'; - if (await persistSnapshotUnavailableProjection( - resources.packageRoots.map((root) => [root, warning] as const), - )) { + const persistableRoots: Array = + error instanceof PiPackageSnapshotLimitError && error.reason === 'duration' + ? [] + : resources.packageRoots.map((root) => [root, { + warning, + durable: error instanceof PiPackageSnapshotLimitError + && error.scope === 'package', + }] as const); + if (await persistSnapshotUnavailableProjection(persistableRoots, inspected)) { await publishPiPackagesChanged({ invalidateCache: false }); } throw error; + } finally { + recordPiPackageStartupDuration(startupTiming, 'package-snapshot', snapshotStartedAt); } }; if (options) return await enqueueMutation(() => resolveResources(true)); await mutationTail; return await resolveResources(); } catch (error) { + if (startupTiming) startupTiming.degraded = true; log.warn('Pi package resources unavailable; starting without user packages', { message: error instanceof Error ? error.message : String(error), }); return { extensions: [], skills: [], promptTemplates: [], packageRoots: [] }; + } finally { + emitPiPackageStartupTiming(startupTiming); } } @@ -1932,12 +2387,33 @@ function enqueueMutation( } class PiPackageSnapshotLimitError extends Error { - constructor(readonly scope: 'package' | 'aggregate' = 'package') { + constructor( + readonly scope: 'package' | 'aggregate' = 'package', + readonly reason: PiPackageSnapshotLimitReason = 'entries', + ) { super('Pi extension snapshot exceeds the safe resource limit'); this.name = 'PiPackageSnapshotLimitError'; } } +function isDurableInspectionFailure(error: unknown): boolean { + return error instanceof PiPackageInspectionLimitError + ? error.reason !== 'duration' + : error instanceof PiPackageSnapshotLimitError + ? error.scope === 'package' && error.reason !== 'duration' + : false; +} + +/** Narrow seam for the inspection-limit persistence decision table. */ +export const __testing = { + isDurableSnapshotLimit( + scope: 'package' | 'aggregate', + reason: PiPackageSnapshotLimitReason, + ): boolean { + return isDurableInspectionFailure(new PiPackageSnapshotLimitError(scope, reason)); + }, +}; + interface SnapshotBudgetCounters { startedAt: number; entries: number; @@ -1972,21 +2448,26 @@ function createSnapshotCopyBudget( }; } -function snapshotBudgetExceeded( +function snapshotBudgetExceededReason( budget: SnapshotBudgetCounters, additionalBytes = 0, -): boolean { - return budget.entries > budget.limits.maxEntries - || budget.bytes + additionalBytes > budget.limits.maxBytes - || Date.now() - budget.startedAt >= budget.limits.maxDurationMs; +): PiPackageSnapshotLimitReason | undefined { + if (budget.entries > budget.limits.maxEntries) return 'entries'; + if (budget.bytes + additionalBytes > budget.limits.maxBytes) return 'bytes'; + if (Date.now() - budget.startedAt >= budget.limits.maxDurationMs) return 'duration'; + return undefined; } function assertSnapshotBudget(budget: SnapshotCopyBudget, additionalBytes = 0): void { - if (snapshotBudgetExceeded(budget, additionalBytes)) { - throw new PiPackageSnapshotLimitError('package'); + const packageReason = snapshotBudgetExceededReason(budget, additionalBytes); + if (packageReason) { + throw new PiPackageSnapshotLimitError('package', packageReason); } - if (budget.aggregate && snapshotBudgetExceeded(budget.aggregate, additionalBytes)) { - throw new PiPackageSnapshotLimitError('aggregate'); + const aggregateReason = budget.aggregate + ? snapshotBudgetExceededReason(budget.aggregate, additionalBytes) + : undefined; + if (aggregateReason) { + throw new PiPackageSnapshotLimitError('aggregate', aggregateReason); } } @@ -2264,9 +2745,14 @@ function mapSnapshotPathOrSkip( : owner.target; } +/** Provenance needed to map staged and skipped roots back to package state. */ interface SnapshotStageMetadata { sourcePackageRoots: string[]; skippedPackageRoots: string[]; + /** Load-dependent duration failures that must remain retryable. */ + transientSkippedPackageRoots: string[]; + /** Package-scoped deterministic limits that may be cached across sessions. */ + durableSkippedPackageRoots: string[]; } const snapshotStageMetadata = new WeakMap(); @@ -2280,15 +2766,18 @@ export async function stageManagedPackageSnapshot( const temporaryRoot = `${snapshotRoot}.tmp-${process.pid}-${Date.now()}`; const mappings: Array<{ source: string; target: string; directory: boolean }> = []; const skippedPackageRoots: string[] = []; + const transientSkippedPackageRoots: string[] = []; + const durableSkippedPackageRoots: string[] = []; const aggregateBudget = createSnapshotBudgetCounters(limits); let aggregateLimitReached = false; + let aggregateLimitReason: PiPackageSnapshotLimitReason | undefined; try { await fs.mkdir(temporaryRoot, { recursive: true, mode: 0o700 }); for (const [index, rawRoot] of resources.packageRoots.entries()) { if (aggregateLimitReached) { - skippedPackageRoots.push( - await fs.realpath(rawRoot).catch(() => path.resolve(rawRoot)), - ); + const skippedRoot = await fs.realpath(rawRoot).catch(() => path.resolve(rawRoot)); + skippedPackageRoots.push(skippedRoot); + if (aggregateLimitReason === 'duration') transientSkippedPackageRoots.push(skippedRoot); continue; } let source: string | undefined; @@ -2317,7 +2806,12 @@ export async function stageManagedPackageSnapshot( mappings.push({ source, target: path.join(snapshotRoot, relativeTarget), directory }); } catch (error) { if (!(error instanceof PiPackageSnapshotLimitError)) throw error; - skippedPackageRoots.push(source ?? path.resolve(rawRoot)); + const skippedRoot = source ?? path.resolve(rawRoot); + skippedPackageRoots.push(skippedRoot); + if (error.reason === 'duration') transientSkippedPackageRoots.push(skippedRoot); + if (error.scope === 'package' && error.reason !== 'duration') { + durableSkippedPackageRoots.push(skippedRoot); + } await fs.rm(path.join(temporaryRoot, String(index)), { recursive: true, force: true, @@ -2325,7 +2819,10 @@ export async function stageManagedPackageSnapshot( // A package-scoped failure quarantines only this root. Existing // mappings remain valid; only the shared aggregate limit stops later // packages from being attempted. - if (error.scope === 'aggregate') aggregateLimitReached = true; + if (error.scope === 'aggregate') { + aggregateLimitReached = true; + aggregateLimitReason = error.reason; + } } } // Windows temp paths can use an 8.3/user-profile spelling while realpath @@ -2354,6 +2851,8 @@ export async function stageManagedPackageSnapshot( snapshotStageMetadata.set(mappedResources, { sourcePackageRoots: mappings.map((mapping) => mapping.source), skippedPackageRoots, + transientSkippedPackageRoots, + durableSkippedPackageRoots, }); return mappedResources; } catch (error) { @@ -2435,7 +2934,7 @@ async function persistEnabledExtensionApprovals(options: { approvedExtensionFingerprints: Object.fromEntries( Object.entries(fingerprints).sort(([left], [right]) => left.localeCompare(right)), ), - snapshotUnavailableRoots: state.snapshotUnavailableRoots, + snapshotUnavailablePackages: state.snapshotUnavailablePackages, }); } @@ -2501,14 +3000,18 @@ async function fingerprintExtensionClosure( const roots = new Set(); let metadataBytes = 0; while (pending.length > 0) { - if (roots.size >= MAX_INSPECTION_ENTRIES) throw new PiPackageInspectionLimitError(); + if (roots.size >= MAX_INSPECTION_ENTRIES) { + throw new PiPackageInspectionLimitError('entries'); + } const root = await fs.realpath(pending.shift()!); if (roots.has(root)) continue; roots.add(root); const manifestResult = await readUtf8FileBounded( path.join(root, 'package.json'), MAX_PACKAGE_JSON_BYTES, root); metadataBytes += manifestResult.bytes; - if (metadataBytes > MAX_INSPECTION_METADATA_BYTES) throw new PiPackageInspectionLimitError(); + if (metadataBytes > MAX_INSPECTION_METADATA_BYTES) { + throw new PiPackageInspectionLimitError('bytes'); + } const manifest = JSON.parse(manifestResult.text) as PackageManifest; const dependencyNames = new Set(Object.keys({ ...manifest.peerDependencies, ...manifest.optionalDependencies, ...manifest.dependencies })); @@ -2592,12 +3095,17 @@ export async function mutatePiPackage( // the existing state cannot be read. Otherwise a transient read failure // could erase explicit disables after a successful package command. await requireState(); + if (await persistSnapshotUnavailableProjection([])) { + mutationMayHaveChangedState = true; + } // Every mutation starts from one fresh projection acquired after the // shared cross-process lock. A packaged/dev/--passive peer may have // installed, removed, updated, or changed approval state since this // process populated its cache; no mutation may persist decisions derived // from that lock-external snapshot. - const inspectedBeforeMutation = await inspectAllPackagesFreshUnderMutationLock(); + const inspectedBeforeMutation = await inspectAllPackagesFreshUnderMutationLock({ + ignorePersistedSnapshotUnavailable: true, + }); const preMutationFreshApprovals = await captureExtensionApprovalIdentities( inspectedBeforeMutation.filter((pkg) => ( pkg.view.requiresExtensionApproval !== true @@ -2651,7 +3159,7 @@ export async function mutatePiPackage( Object.entries(state.approvedExtensionFingerprints) .filter(([item]) => !removedSources.has(item)), ), - snapshotUnavailableRoots: state.snapshotUnavailableRoots, + snapshotUnavailablePackages: state.snapshotUnavailablePackages, }); } else if (request.action === 'update') { mutationMayHaveChangedState = true; @@ -2729,7 +3237,7 @@ export async function mutatePiPackage( approvedExtensionFingerprints: Object.fromEntries( Object.entries(approvedFingerprints).sort(([left], [right]) => left.localeCompare(right)), ), - snapshotUnavailableRoots: state.snapshotUnavailableRoots, + snapshotUnavailablePackages: state.snapshotUnavailablePackages, }); } invalidateInspectionCache(); diff --git a/packages/maker-core/src/agents/base-agent.ts b/packages/maker-core/src/agents/base-agent.ts index 7191b3a7cc..5607c7f60b 100644 --- a/packages/maker-core/src/agents/base-agent.ts +++ b/packages/maker-core/src/agents/base-agent.ts @@ -550,7 +550,11 @@ export interface AgentDeps { * and path confinement. Device-link remote control still executes on this host * and therefore uses these resources. SSH remoteHostId and Review runtimes do not. */ - resolvePiManagedPackageResources?: (options?: { snapshotRoot: string }) => Promise<{ + resolvePiManagedPackageResources?: (options?: { + snapshotRoot: string; + /** Redacted per-start correlation id for structured startup timing logs. */ + startupTraceId?: string; + }) => Promise<{ extensions: string[]; skills: Array<{ path: string; name: string; description?: string }>; promptTemplates: string[]; diff --git a/packages/maker-core/src/agents/pi/__tests__/pi-agent.integration.test.ts b/packages/maker-core/src/agents/pi/__tests__/pi-agent.integration.test.ts index cb7ae72883..e4c4b607c4 100644 --- a/packages/maker-core/src/agents/pi/__tests__/pi-agent.integration.test.ts +++ b/packages/maker-core/src/agents/pi/__tests__/pi-agent.integration.test.ts @@ -2219,6 +2219,7 @@ describe.skipIf(!piAvailable)('PiAgent integration (real pi binary + fake gatewa symlinkSync('../secrets/.env', postCdLink); symlinkSync(ordinaryPath, path.join(workingDir, 'link')); symlinkSync(ordinaryPath, path.join(stackOtherDir, 'link')); + mkdirSync(path.dirname(path.join(workingDir, escapedLinkName)), { recursive: true }); symlinkSync(secretPath, path.join(workingDir, escapedLinkName)); symlinkSync(secretPath, path.join(workingDir, cdRedirectLinkName)); symlinkSync(ordinaryPath, path.join(subDir, cdRedirectLinkName)); diff --git a/packages/maker-core/src/agents/pi/__tests__/pi-provider-routing.test.ts b/packages/maker-core/src/agents/pi/__tests__/pi-provider-routing.test.ts index c240abf76c..2701682819 100644 --- a/packages/maker-core/src/agents/pi/__tests__/pi-provider-routing.test.ts +++ b/packages/maker-core/src/agents/pi/__tests__/pi-provider-routing.test.ts @@ -18,6 +18,10 @@ const captured = vi.hoisted(() => ({ refreshTimeoutOnEvent?: (event: { type: string }) => boolean; } | undefined>, closes: 0, + onExit: undefined as undefined | ((info: { + code: number | null; + signal: NodeJS.Signals | null; + }) => void), requestHandler: undefined as undefined | ((command: Record) => Promise<{ success: boolean; command?: unknown; @@ -64,6 +68,14 @@ vi.mock('../rpc-client.js', () => { PiRpcRequestTimeoutError, PiRpcProcess: class { isClosed = false; + constructor(opts: { + onExit: (info: { + code: number | null; + signal: NodeJS.Signals | null; + }) => void; + }) { + captured.onExit = opts.onExit; + } async request( command: Record, options?: { @@ -129,6 +141,7 @@ describe('Pi provider-aware model routing', () => { captured.requests = []; captured.requestOptions = []; captured.closes = 0; + captured.onExit = undefined; captured.requestHandler = undefined; agentHome = mkdtempSync(path.join(tmpdir(), 'pi-provider-home-')); cwd = mkdtempSync(path.join(tmpdir(), 'pi-provider-cwd-')); @@ -3705,6 +3718,66 @@ describe('Pi provider-aware model routing', () => { await handle.close(); }); + it('preserves the stable remote config home across a transport disconnect and reattach', async () => { + const remoteStub: import('../transport.js').PiTransport = { + writeLine: async () => {}, + onLine: () => () => {}, + onStderr: () => () => {}, + onClose: () => () => {}, + close: async () => {}, + pid: 4321, + isClosed: () => false, + remoteBinaryPath: '/remote/pi', + killRemoteSession: async () => {}, + }; + const remoteRm = vi.fn(async () => {}); + const capturedRemoteEnvs: Array> = []; + const base = byomDeps(async () => ({ providers: [], env: {} })); + const deps: AgentDeps = { + ...base, + runtimeConfig: { + ...base.runtimeConfig, + remoteEndpoint: 'https://gateway.example.test', + }, + resolveRemotePiBinaryPath: async () => '/remote/pi', + getRemotePiTransport: async (_hostId, opts) => { + capturedRemoteEnvs.push({ ...(opts.env ?? {}) }); + return remoteStub; + }, + getRemotePiFileOps: () => ({ + mkdirp: async () => {}, + writeFile: async () => {}, + stat: async () => ({ isFile: true }), + rm: remoteRm, + listDir: async () => [], + }), + }; + + await new PiAgent(deps).startSession({ + sessionId: 'remote-disconnect-reattach', + workingDir: cwd, + model: 'local-model', + remoteHostId: 'remote-host', + }); + const firstEnv = capturedRemoteEnvs[0]!; + const configHome = firstEnv.PI_CODING_AGENT_DIR!; + + captured.onExit?.({ code: null, signal: null }); + await Promise.resolve(); + expect(remoteRm).not.toHaveBeenCalledWith(configHome, { recursive: true }); + + const reattached = await new PiAgent(deps).startSession({ + sessionId: 'remote-disconnect-reattach', + workingDir: cwd, + model: 'local-model', + remoteHostId: 'remote-host', + }); + expect(capturedRemoteEnvs[1]).toEqual(firstEnv); + expect(capturedRemoteEnvs[1]?.PI_CODING_AGENT_DIR).toBe(configHome); + await reattached.close(); + expect(remoteRm).not.toHaveBeenCalledWith(configHome, { recursive: true }); + }); + it('hashes the remote permission snapshot into spawn env so a later Full-access attach restarts', async () => { const remoteStub: import('../transport.js').PiTransport = { writeLine: async () => {}, diff --git a/packages/maker-core/src/agents/pi/__tests__/pi-startsession-cleanup.test.ts b/packages/maker-core/src/agents/pi/__tests__/pi-startsession-cleanup.test.ts index 9f739fbec5..9474ec5be9 100644 --- a/packages/maker-core/src/agents/pi/__tests__/pi-startsession-cleanup.test.ts +++ b/packages/maker-core/src/agents/pi/__tests__/pi-startsession-cleanup.test.ts @@ -13,6 +13,7 @@ */ import { existsSync, mkdirSync, mkdtempSync, readFileSync, realpathSync, rmSync, writeFileSync } from 'node:fs'; +import { createHash } from 'node:crypto'; import { tmpdir } from 'node:os'; import path from 'node:path'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; @@ -2187,9 +2188,10 @@ describe('PiAgent.startSession failure cleanup (mocked pi process)', () => { sessionId: 'approved', })); expect(resolvePiManagedPackageResources).toHaveBeenCalledOnce(); - expect(resolvePiManagedPackageResources).toHaveBeenCalledWith({ + expect(resolvePiManagedPackageResources).toHaveBeenCalledWith(expect.objectContaining({ snapshotRoot: path.join(approvedHome, 'managed-packages'), - }); + startupTraceId: expect.stringMatching(/^[a-f0-9]{16}$/), + })); await vi.waitFor(() => { expect(approvedHandle.getRuntimeCapabilities?.()?.projectResources).toMatchObject({ status: 'approved', approvalRevision: 'rev-approved-only', requestedSkillCount: 1, @@ -2219,9 +2221,310 @@ describe('PiAgent.startSession failure cleanup (mocked pi process)', () => { // 上层随后仍可能调用 close() —— cleanup 幂等,不抛。 await handle.close(); }); + + it("registers ownership before package setup and cleans the config home when setup fails", async () => { + const { promises: fs } = await import("node:fs"); + const originalWriteFile = fs.writeFile.bind(fs); + let configHome = ""; + let ownerWasPresent = false; + const writeSpy = vi + .spyOn(fs, "writeFile") + .mockImplementation(async (target, ...args) => { + if (path.basename(String(target)) === "models.json") { + configHome = path.dirname(String(target)); + ownerWasPresent = existsSync( + path.join(configHome, ".cindy-owner.json"), + ); + throw new Error("models setup failed (mock)"); + } + return originalWriteFile( + target, + ...(args as Parameters extends [ + unknown, + ...infer Rest, + ] + ? Rest + : never), + ); + }); + try { + const agent = new PiAgent(buildDeps()); + await expect( + agent.startSession({ + sessionId: "setup-failure", + workingDir: cwd, + model: "m", + }), + ).rejects.toThrow("models setup failed (mock)"); + expect(ownerWasPresent).toBe(true); + expect(configHome).not.toBe(""); + await waitFor(() => !existsSync(configHome)); + } finally { + writeSpy.mockRestore(); + } + }); + + it('records process incarnation plus redacted session/runtime identity in the local owner marker', async () => { + const sessionId = 'owner-marker-session'; + const handle = await new PiAgent(buildDeps()).startSession({ + sessionId, + workingDir: cwd, + model: 'm', + }); + const configHome = knobs.spawnedEnvs[0]!.PI_CODING_AGENT_DIR!; + try { + const owner = JSON.parse( + readFileSync(path.join(configHome, '.cindy-owner.json'), 'utf8'), + ) as Record; + expect(owner).toMatchObject({ + version: 2, + ownerPid: process.pid, + directoryName: path.basename(configHome), + runtimeId: path.basename(configHome), + sessionIdHash: createHash('sha256').update(sessionId).digest('hex'), + }); + expect(owner.ownerStartTimeSec).toEqual(expect.any(Number)); + expect(owner).not.toHaveProperty('sessionId'); + } finally { + await handle.close(); + } + }); + + it('correlates redacted MCP, spawn, RPC-ready, and first-model-request timings', async () => { + const info = vi.fn(); + const timingLogger: Logger = { + ...noopLogger, + info, + child: () => timingLogger, + }; + const resolvePiManagedPackageResources = vi.fn(async () => ({ + extensions: [], skills: [], promptTemplates: [], packageRoots: [], + })); + const handle = await new PiAgent(buildDeps({ + logger: timingLogger, + resolvePiManagedPackageResources, + })).startSession({ + sessionId: 'timing-session-secret', + workingDir: cwd, + model: 'm', + }); + try { + await handle.send({ type: 'user', content: 'timing-prompt-secret' }); + const resolverOptions = resolvePiManagedPackageResources.mock.calls[0]?.[0] as + | { startupTraceId?: string } + | undefined; + expect(resolverOptions?.startupTraceId).toMatch(/^[a-f0-9]{16}$/); + + const events = info.mock.calls + .filter(([message]) => message === 'pi startup stage') + .map(([, fields]) => fields as Record); + expect(events.map((event) => event.stage)).toEqual([ + 'mcp-ready', + 'pi-spawn', + 'rpc-ready', + 'first-model-request', + ]); + for (const event of events) { + expect(event).toMatchObject({ + startupTraceId: resolverOptions?.startupTraceId, + stage: expect.any(String), + durationMs: expect.any(Number), + status: 'ok', + }); + } + const serialized = JSON.stringify(events); + expect(serialized).not.toContain('timing-session-secret'); + expect(serialized).not.toContain('timing-prompt-secret'); + expect(serialized).not.toContain(cwd); + } finally { + await handle.close(); + } + }); + + it('reclaims a v2 local config home whose owner pid was recycled', async () => { + const runtimeId = 'a'.repeat(32); + const recycledHome = path.join(agentHome, 'run-tmp', runtimeId); + mkdirSync(recycledHome, { recursive: true }); + writeFileSync( + path.join(recycledHome, '.cindy-owner.json'), + `${JSON.stringify({ + version: 2, + ownerPid: process.pid, + ownerStartTimeSec: 1, + createdAt: 1, + directoryName: path.basename(recycledHome), + runtimeId, + sessionIdHash: 'b'.repeat(64), + })}\n`, + ); + + const handle = await new PiAgent(buildDeps()).startSession({ + sessionId: 'after-pid-reuse', + workingDir: cwd, + model: 'm', + }); + try { + expect(existsSync(recycledHome)).toBe(false); + } finally { + await handle.close(); + } + }); + + it('shares one process-start probe memo across config homes owned by the same process', async () => { + const ownerPid = process.ppid; + const ownerStartTimeSec = 1; + const ownerHomes = ['c'.repeat(32), 'd'.repeat(32)].map((runtimeId) => { + const ownerHome = path.join(agentHome, 'run-tmp', runtimeId); + mkdirSync(ownerHome, { recursive: true }); + writeFileSync( + path.join(ownerHome, '.cindy-owner.json'), + `${JSON.stringify({ + version: 2, + ownerPid, + ownerStartTimeSec, + createdAt: 1, + directoryName: runtimeId, + runtimeId, + sessionIdHash: 'e'.repeat(64), + })}\n`, + ); + return ownerHome; + }); + const probeMemos: unknown[] = []; + const ownerProbe = vi + .spyOn(piSubagentRuns, 'isPiHostProcessInstanceAlive') + .mockImplementation(function () { + probeMemos.push(arguments[1]); + return true; + }); + + const handle = await new PiAgent(buildDeps()).startSession({ + sessionId: 'shared-owner-probe-memo', + workingDir: cwd, + model: 'm', + }); + try { + expect(ownerProbe).toHaveBeenCalledTimes(ownerHomes.length); + // pi-subagent-runs separately locks that one shared memo means one + // PowerShell/ps start-time probe per owner pid and sweep. + expect(probeMemos[0]).toBeInstanceOf(Map); + expect(probeMemos[1]).toBe(probeMemos[0]); + for (const ownerHome of ownerHomes) expect(existsSync(ownerHome)).toBe(true); + } finally { + await handle.close(); + } + }); + + it('preserves markerless legacy homes even when the host process scan finds no live Pi', async () => { + const legacyHome = path.join(agentHome, 'run-tmp', '0123456789abcdef'); + mkdirSync(legacyHome, { recursive: true }); + writeFileSync(path.join(legacyHome, 'models.json'), '{}\n'); + + // Simulate the removed host hint so a future sweep cannot accidentally + // reinstate markerless deletion based only on post-spawn process evidence. + const deps = { + ...buildDeps(), + canReclaimLegacyPiConfigHomes: async () => true, + }; + const handle = await new PiAgent(deps).startSession({ + sessionId: 'legacy-preserve-pre-spawn', + workingDir: cwd, + model: 'm', + }); + try { + expect(existsSync(legacyHome)).toBe(true); + } finally { + await handle.close(); + } + }); + + it("reclaims failed cleanup without deleting same-process or other-live-process sessions", async () => { + const { promises: fs } = await import("node:fs"); + const agent = new PiAgent(buildDeps()); + const orphaned = await agent.startSession({ + sessionId: "orphaned", + workingDir: cwd, + model: "m", + }); + const active = await agent.startSession({ + sessionId: "active", + workingDir: cwd, + model: "m", + }); + const orphanedHome = knobs.spawnedEnvs.find( + (env) => env.CINDY_PI_SESSION_ID === "orphaned", + )!.PI_CODING_AGENT_DIR!; + const activeHome = knobs.spawnedEnvs.find( + (env) => env.CINDY_PI_SESSION_ID === "active", + )!.PI_CODING_AGENT_DIR!; + const liveOwnerHome = path.join(agentHome, "run-tmp", "live-owner-home"); + const originalRm = fs.rm.bind(fs); + let orphanRemovalFailures = 3; + const rmSpy = vi + .spyOn(fs, "rm") + .mockImplementation(async (target, options) => { + if ( + path.resolve(String(target)) === path.resolve(orphanedHome) && + orphanRemovalFailures > 0 + ) { + orphanRemovalFailures -= 1; + throw Object.assign(new Error("transient config-home lock"), { + code: "EPERM", + }); + } + return originalRm(target, options); + }); + try { + await orphaned.close(); + await waitFor(() => orphanRemovalFailures === 0); + expect(existsSync(orphanedHome)).toBe(true); + expect(existsSync(path.join(activeHome, "models.json"))).toBe(true); + + const deadOwnerHome = path.join(agentHome, "run-tmp", "dead-owner-home"); + mkdirSync(deadOwnerHome, { recursive: true }); + writeFileSync( + path.join(deadOwnerHome, ".cindy-owner.json"), + `${JSON.stringify({ + version: 1, + ownerPid: 2_147_483_647, + createdAt: 1, + directoryName: path.basename(deadOwnerHome), + })}\n`, + ); + expect(process.ppid).toBeGreaterThan(0); + mkdirSync(liveOwnerHome, { recursive: true }); + writeFileSync( + path.join(liveOwnerHome, ".cindy-owner.json"), + `${JSON.stringify({ + version: 1, + ownerPid: process.ppid, + createdAt: 1, + directoryName: path.basename(liveOwnerHome), + })}\n`, + ); + const next = await agent.startSession({ + sessionId: "next", + workingDir: cwd, + model: "m", + }); + try { + expect(existsSync(orphanedHome)).toBe(false); + expect(existsSync(deadOwnerHome)).toBe(false); + expect(existsSync(path.join(activeHome, "models.json"))).toBe(true); + expect(existsSync(liveOwnerHome)).toBe(true); + } finally { + await next.close(); + } + } finally { + rmSpy.mockRestore(); + await active.close(); + await originalRm(orphanedHome, { recursive: true, force: true }); + await originalRm(liveOwnerHome, { recursive: true, force: true }); + } + }); }); -/** 轮询等待条件成立(configHome cleanup 是 void fs.rm fire-and-forget,不阻塞 close)。 */ +/** 轮询等待条件成立(onExit 等同步回调里的 configHome cleanup 仍是 fire-and-forget)。 */ async function waitFor(cond: () => boolean, timeoutMs = 2000): Promise { const start = Date.now(); while (!cond()) { diff --git a/packages/maker-core/src/agents/pi/index.ts b/packages/maker-core/src/agents/pi/index.ts index 30ef1e84a4..e0a0649dfb 100644 --- a/packages/maker-core/src/agents/pi/index.ts +++ b/packages/maker-core/src/agents/pi/index.ts @@ -100,6 +100,8 @@ import { listPiSubagentRunDiagnostics, listPiSubagentRunDirectoryIds, listPiSubagentRuns, + isPiHostProcessInstanceAlive, + piHostProcessStartTimeSec, piSubagentRunRoot, piSubagentApprovalScope, piSubagentRuntimeOwnerId, @@ -111,6 +113,7 @@ import { syncPiSubagentPermissions, type PiSubagentRunDiagnostic, type PiSubagentRunStatus, + type ProcessStartTimeMemo, } from './pi-subagent-runs.js'; import { annotatePermissionRequestForUnavailableReview, @@ -449,6 +452,248 @@ async function stageManagedRipgrep(configHome: string, sourcePath: string | unde return targetPath; } +const LOCAL_CONFIG_HOME_OWNER_FILE = '.cindy-owner.json'; +const LOCAL_CONFIG_HOME_OWNER_VERSION = 2; +const LOCAL_CONFIG_HOME_REMOVE_RETRY_DELAYS_MS = [0, 25, 100] as const; +const LOCAL_CONFIG_HOME_RUNTIME_ID_RE = /^[a-f0-9]{32}$/; +const SHA256_HEX_RE = /^[a-f0-9]{64}$/; +const activeLocalConfigHomes = new Set(); + +/** Durable proof used only to decide whether a local per-session config home is reclaimable. */ +interface LocalConfigHomeOwnerV1 { + version: 1; + ownerPid: number; + createdAt: number; + directoryName: string; +} + +interface LocalConfigHomeOwnerV2 { + version: typeof LOCAL_CONFIG_HOME_OWNER_VERSION; + ownerPid: number; + ownerStartTimeSec: number; + createdAt: number; + directoryName: string; + runtimeId: string; + sessionIdHash: string; +} + +type LocalConfigHomeOwner = LocalConfigHomeOwnerV1 | LocalConfigHomeOwnerV2; + +/** Minimal logging interface required by local config-home reclamation policy. */ +interface ConfigHomeCleanupLogger { + debug(message: string, fields?: Record): void; +} + +type PiStartupStage = 'mcp-ready' | 'pi-spawn' | 'rpc-ready' | 'first-model-request'; + +function logPiStartupStage( + logger: AgentDeps['logger'], + startupTraceId: string, + stage: PiStartupStage, + startedAt: number, + status: 'ok' | 'degraded' | 'skipped', + counts: Record = {}, +): void { + logger.info('pi startup stage', { + startupTraceId, + stage, + durationMs: Math.max(0, Date.now() - startedAt), + status, + ...counts, + }); +} + +function localConfigHomeKey(configHome: string): string { + const resolved = path.resolve(configHome); + return process.platform === 'win32' ? resolved.toLowerCase() : resolved; +} + +async function registerLocalConfigHome( + configHome: string, + sessionId: string, + runtimeId: string, +): Promise { + const normalized = path.resolve(configHome); + if (path.basename(normalized) !== runtimeId || !LOCAL_CONFIG_HOME_RUNTIME_ID_RE.test(runtimeId)) { + throw new Error('pi: invalid local config-home runtime identity'); + } + const owner: LocalConfigHomeOwnerV2 = { + version: LOCAL_CONFIG_HOME_OWNER_VERSION, + ownerPid: process.pid, + ownerStartTimeSec: piHostProcessStartTimeSec(), + createdAt: Date.now(), + directoryName: path.basename(normalized), + runtimeId, + sessionIdHash: createHash('sha256').update(sessionId).digest('hex'), + }; + await fs.writeFile( + path.join(normalized, LOCAL_CONFIG_HOME_OWNER_FILE), + `${JSON.stringify(owner)}\n`, + { encoding: 'utf8', mode: 0o600, flag: 'wx' }, + ); + activeLocalConfigHomes.add(localConfigHomeKey(normalized)); +} + +function unregisterLocalConfigHome(configHome: string): void { + activeLocalConfigHomes.delete(localConfigHomeKey(configHome)); +} + +async function readLocalConfigHomeOwner(configHome: string): Promise { + const markerPath = path.join(configHome, LOCAL_CONFIG_HOME_OWNER_FILE); + try { + const markerStat = await fs.lstat(markerPath); + if (!markerStat.isFile() || markerStat.size > 1024) return null; + const parsed = JSON.parse(await fs.readFile(markerPath, 'utf8')) as Record; + if ( + !Number.isSafeInteger(parsed.ownerPid) + || (parsed.ownerPid as number) <= 0 + || !Number.isSafeInteger(parsed.createdAt) + || (parsed.createdAt as number) <= 0 + || parsed.directoryName !== path.basename(configHome) + ) return null; + if (parsed.version === 1) { + return { + version: 1, + ownerPid: parsed.ownerPid as number, + createdAt: parsed.createdAt as number, + directoryName: parsed.directoryName, + }; + } + if ( + parsed.version !== LOCAL_CONFIG_HOME_OWNER_VERSION + || !Number.isSafeInteger(parsed.ownerStartTimeSec) + || (parsed.ownerStartTimeSec as number) <= 0 + || typeof parsed.runtimeId !== 'string' + || parsed.runtimeId !== parsed.directoryName + || !LOCAL_CONFIG_HOME_RUNTIME_ID_RE.test(parsed.runtimeId) + || typeof parsed.sessionIdHash !== 'string' + || !SHA256_HEX_RE.test(parsed.sessionIdHash) + ) return null; + return { + version: LOCAL_CONFIG_HOME_OWNER_VERSION, + ownerPid: parsed.ownerPid as number, + ownerStartTimeSec: parsed.ownerStartTimeSec as number, + createdAt: parsed.createdAt as number, + directoryName: parsed.directoryName, + runtimeId: parsed.runtimeId, + sessionIdHash: parsed.sessionIdHash, + }; + } catch { + return null; + } +} + +function isLocalProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch (error) { + return (error as NodeJS.ErrnoException).code !== 'ESRCH'; + } +} + +function localConfigHomeOwnerIsActive( + owner: LocalConfigHomeOwner, + configHome: string, + startTimeMemo?: ProcessStartTimeMemo, +): boolean { + if (owner.version === 1) { + return owner.ownerPid === process.pid + ? activeLocalConfigHomes.has(localConfigHomeKey(configHome)) + : isLocalProcessAlive(owner.ownerPid); + } + if (!isPiHostProcessInstanceAlive({ + pid: owner.ownerPid, + startTimeSec: owner.ownerStartTimeSec, + }, startTimeMemo)) return false; + return owner.ownerPid !== process.pid + || activeLocalConfigHomes.has(localConfigHomeKey(configHome)); +} + +function sameLocalConfigHomeOwner( + left: LocalConfigHomeOwner, + right: LocalConfigHomeOwner, +): boolean { + if (left.version !== right.version) return false; + if ( + left.ownerPid !== right.ownerPid + || left.createdAt !== right.createdAt + || left.directoryName !== right.directoryName + ) return false; + return left.version === 1 || ( + right.version === LOCAL_CONFIG_HOME_OWNER_VERSION + && left.ownerStartTimeSec === right.ownerStartTimeSec + && left.runtimeId === right.runtimeId + && left.sessionIdHash === right.sessionIdHash + ); +} + +async function removeLocalConfigHomeWithRetry( + configHome: string, +): Promise<{ removed: true } | { removed: false; error: unknown }> { + let lastError: unknown = new Error('pi configHome cleanup failed'); + for (const delayMs of LOCAL_CONFIG_HOME_REMOVE_RETRY_DELAYS_MS) { + if (delayMs > 0) { + await new Promise((resolve) => setTimeout(resolve, delayMs)); + } + try { + await fs.rm(configHome, { recursive: true, force: true }); + return { removed: true }; + } catch (error) { + lastError = error; + } + } + return { removed: false, error: lastError }; +} + +async function sweepStaleLocalConfigHomes( + agentHome: string, + logger: ConfigHomeCleanupLogger, +): Promise { + const runTmp = path.resolve(agentHome, 'run-tmp'); + let entries; + try { + entries = await fs.readdir(runTmp, { withFileTypes: true }); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return; + logger.debug('pi configHome orphan scan failed (non-fatal)', { + message: error instanceof Error ? error.message : String(error), + }); + return; + } + // One owner process may have many active session homes. Reuse only the + // expensive start-time sample within this sweep; each check still performs + // a fresh signal-0 liveness probe and deletion still re-reads the marker. + const ownerStartTimeMemo: ProcessStartTimeMemo = new Map(); + for (const entry of entries) { + if (!entry.isDirectory()) continue; + const candidate = path.resolve(runTmp, entry.name); + if (path.dirname(candidate) !== runTmp) continue; + const owner = await readLocalConfigHomeOwner(candidate); + // Markerless homes may belong to an older Cindy instance between directory + // creation and Pi spawn. No process snapshot can prove them reclaimable. + if (!owner) continue; + if (localConfigHomeOwnerIsActive(owner, candidate, ownerStartTimeMemo)) continue; + + // Re-read immediately before deletion. A directory whose marker changed or + // whose owner became active is no longer the orphan we proved above. + const currentOwner = await readLocalConfigHomeOwner(candidate); + if ( + !currentOwner + || !sameLocalConfigHomeOwner(currentOwner, owner) + ) continue; + if (localConfigHomeOwnerIsActive(currentOwner, candidate, ownerStartTimeMemo)) continue; + + const outcome = await removeLocalConfigHomeWithRetry(candidate); + if (!outcome.removed) { + logger.debug('pi configHome orphan cleanup failed (retried next startSession)', { + configHomeId: entry.name, + message: outcome.error instanceof Error ? outcome.error.message : String(outcome.error), + }); + } + } +} + /** cindy Effort → pi thinking level(pi 无 ultra)。思考开关走 setThinkingEnabled / thinkingEnabled。 */ function effortToPiThinkingLevel(effort: Effort): string { return effort === 'ultra' ? 'max' : effort; @@ -1596,6 +1841,7 @@ export class PiAgent extends BaseAgent { } private async startSessionWhileRunning(opts: StartSessionOptions): Promise { + const startupTraceId = randomBytes(8).toString('hex'); const startupCleanupKey = opts.sessionId ?? ''; // A previous pre-publication Pi process for this business session must be // confirmed dead before another spawn can begin. @@ -2039,6 +2285,12 @@ export class PiAgent extends BaseAgent { message: 'remote pi sessions require getRemotePiFileOps — host must provide SSH file primitives', }); } + if (!remote) { + await sweepStaleLocalConfigHomes( + agentHome, + this.deps.logger, + ); + } const allowPiPackageManagement = !reviewMode && !remote && Boolean(this.deps.mutatePiManagedPackage); // The UI-request title is visible to every extension in the Pi process. // Authenticate the host-backed mutation channel with a per-runtime bearer @@ -2068,7 +2320,11 @@ export class PiAgent extends BaseAgent { return null; } }; - const rmPath = async (target: string, opts2?: { recursive?: boolean }): Promise => { + const rmPath = async ( + target: string, + opts2?: { recursive?: boolean }, + propagateError = false, + ): Promise => { // 轮 22 MEDIUM-5:两分支统一吞错(force 语义) —— 远端 fileOps.rm 抛错 // 时本地分支却静默, 语义不对称会误导未来调用方。rm 失败由上层 // cleanupConfigHome 的 catch 兜底。 @@ -2081,7 +2337,8 @@ export class PiAgent extends BaseAgent { force: true, }); } - } catch { + } catch (error) { + if (propagateError) throw error; /* best-effort:rm 失败不致命 */ } }; @@ -2103,36 +2360,58 @@ export class PiAgent extends BaseAgent { // 远端再叠 models.json+settings.json hash:startSession 在 pi/ensure 之前就会写 // 这两份快照。若只按 sessionId 分目录, 另一实例改路由或 retry 策略会先覆盖 // 仍在跑的旧 Pi / 子代理热读快照。 + const localConfigHomeRuntimeId = remote ? undefined : randomBytes(16).toString('hex'); let configHome = remote ? joinRemotePosixPath(agentHome, 'run-tmp', stableSessionPathSegment(opts.sessionId)) - : joinRemotePosixPath(agentHome, 'run-tmp', randomBytes(8).toString('hex')); + : joinRemotePosixPath(agentHome, 'run-tmp', localConfigHomeRuntimeId!); let configHomeCleaned = false; - // 清理失败(SSH 断链时 fileOps.rm 抛错)不置标志 —— 下次会话的 startSession - // 会主动清陈旧 configHome(见下),且不因「一次失败永久跳过」累积泄漏 - // (R4-2 竞态 1/6)。 - const cleanupConfigHome = (): void => { - if (configHomeCleaned) return; - void rmPath(configHome, { recursive: true }).then( - () => { - configHomeCleaned = true; - }, - // 失败不置标志(下次 startSession 清陈旧目录兜底), 但必须留日志—— - // 否则远端 SSH fs 卡死等根因不可见(R7 审计 L-3)。 - (err) => { - this.deps.logger.debug('pi configHome cleanup failed (retried next startSession)', { - configHome, - message: err instanceof Error ? err.message : String(err), - }); - }, - ); + let configHomeCleanupPromise: Promise | undefined; + const cleanupConfigHome = (): Promise => { + if (configHomeCleaned) return Promise.resolve(); + if (configHomeCleanupPromise) return configHomeCleanupPromise; + if (!remote) unregisterLocalConfigHome(configHome); + const pending = (async () => { + if (remote) { + await rmPath(configHome, { recursive: true }, true); + } else { + const outcome = await removeLocalConfigHomeWithRetry(configHome); + if (!outcome.removed) throw outcome.error; + } + configHomeCleaned = true; + })().catch((err) => { + // Keep the owner marker and cleaned=false. A later lifecycle callback + // retries directly; the next local start also reclaims this directory + // after proving that its owner/session is no longer active. + this.deps.logger.debug('pi configHome cleanup failed (retried next startSession)', { + configHomeId: remote ? path.posix.basename(configHome) : path.basename(configHome), + message: err instanceof Error ? err.message : String(err), + }); + }).finally(() => { + if (configHomeCleanupPromise === pending) configHomeCleanupPromise = undefined; + }); + configHomeCleanupPromise = pending; + return pending; }; - // 轮 40-w4-t4 CRITICAL:不再在新会话启动时清 run-tmp 其它目录 —— 远端 - // 并发会话 A 的 configHome 会被 B 的启动清理删除(无 owner/lease 校验), - // 破坏 A 的 bridge extension/models.json 运行期快照。清理只绑定到 - // 本会话自己的 close/失败(cleanupConfigHome);run-tmp 残留(断链时清理 - // 失败的孤儿)由本会话 close 路径的低频清理覆盖, 不牺牲活跃会话。 - // 注:断链残留的孤儿 configHome 会累积 —— 但删错活跃会话是毁任务, - // 宁可残留(可由用户手动清或未来加 lease 机制), 不误删。 + // Until a handle owns the lifecycle, any setup rejection must remove this + // private home. An unconfirmed process close transfers cleanup to the + // session-keyed quarantine instead, because its Pi process may still read it. + let configHomeLifecycleTransferred = false; + let configHomeCleanupDeferredToQuarantine = false; + let localConfigHomeCreated = false; + try { + if (!remote) { + await fs.mkdir(path.dirname(configHome), { recursive: true, mode: 0o700 }); + await fs.mkdir(configHome, { mode: 0o700 }); + localConfigHomeCreated = true; + await registerLocalConfigHome( + configHome, + opts.sessionId ?? '', + localConfigHomeRuntimeId!, + ); + } + // Local run-tmp sweeping is owner-aware: a live pid, or a configHome still + // registered by this process, is never removed. Remote daemon homes keep + // their existing stable identity and are not swept by local filesystem logic. if (remote) { const preview = await this.writeModelsJson(configHome, nativeProviders, retainedRuntimeModel, authProviderId, { remote, @@ -2465,6 +2744,8 @@ export class PiAgent extends BaseAgent { // Phase 1 不桥 orca/memory/ghost)。外部 HTTP MCP 直连不受影响。 const remoteSkipMcpBridge = Boolean( opts.remoteHostId && this.deps.remotePiSkipMcpBridge?.(opts.remoteHostId)); + const mcpReadyStartedAt = Date.now(); + let mcpReadyStatus: 'ok' | 'degraded' | 'skipped' = 'skipped'; if (!reviewMode && !remoteSkipMcpBridge && this.deps.preparePiExtraSpawnConfig) { try { const extra = await this.deps.preparePiExtraSpawnConfig(this.deps.mcpProviders ?? [], { @@ -2486,12 +2767,22 @@ export class PiAgent extends BaseAgent { registeredMcpServerNames.add(server.name); } } + mcpReadyStatus = 'ok'; } catch (err) { + mcpReadyStatus = 'degraded'; this.deps.logger.error('pi MCP bridge prep failed, continuing without cindy tools', { message: err instanceof Error ? err.message : String(err), }); } } + logPiStartupStage( + this.deps.logger, + startupTraceId, + 'mcp-ready', + mcpReadyStartedAt, + mcpReadyStatus, + { mcpServerCount: registeredMcpServerNames.size }, + ); // 压缩即记忆:makerMemory 开启时,把 pi 压缩上下文时丢弃内容的摘要沉淀成 `digest` // 记忆(进 FTS 可 memory_search 检索,但排除出 MEMORY.md / system prompt,不污染 @@ -2624,6 +2915,7 @@ export class PiAgent extends BaseAgent { try { managedPackageResources = await this.deps.resolvePiManagedPackageResources({ snapshotRoot: path.join(configHome, 'managed-packages'), + startupTraceId, }); } catch { this.deps.logger.warn('pi managed package resolver failed closed', { @@ -3802,6 +4094,8 @@ export class PiAgent extends BaseAgent { return proxyLeaseInitialInspection; }; let durableSpawnEnv: NodeJS.ProcessEnv = {}; + let piSpawnStartedAt: number | undefined; + let piSpawnLogged = false; try { // 远端不 stage 本地 rg(本机二进制远端无意义)—— 远端走 PATH 上的 rg(远端 POSIX // 系统常见),与 CC/Codex 远端一致(不注入受管工具路径)。 @@ -3933,6 +4227,7 @@ export class PiAgent extends BaseAgent { mergeLoopbackNoProxy(spawnEnv); durableSpawnEnv = spawnEnv; const initialHostProxyForward = nativeProviderById.get(initialProvider)?.hostProxyForward; + piSpawnStartedAt = Date.now(); const { transport } = await this.createTransport( { args, @@ -4154,13 +4449,18 @@ export class PiAgent extends BaseAgent { // runtime 文件残留由下次 startSession 的清陈旧目录兜底(轮 40-w4-t4)。 // 本地 stdio 无 daemon, onExit 即真死, 保持清理。 if (!remote) { - cleanupConfigHome(); + void cleanupConfigHome(); cleanupRuntimeFiles(); } queue.end(); }, }); + logPiStartupStage(this.deps.logger, startupTraceId, 'pi-spawn', piSpawnStartedAt, 'ok'); + piSpawnLogged = true; } catch (err) { + if (piSpawnStartedAt !== undefined && !piSpawnLogged) { + logPiStartupStage(this.deps.logger, startupTraceId, 'pi-spawn', piSpawnStartedAt, 'degraded'); + } clearPiSubagentRefreshTimer(); disposePiTranslateContext(ctx); try { @@ -4170,7 +4470,7 @@ export class PiAgent extends BaseAgent { } // 轮 42 P1:远端失败也不清理 runtime 文件(可能与并发存活会话共享/复用)。 if (!remote) { - cleanupConfigHome(); + await cleanupConfigHome(); cleanupRuntimeFiles(); } throw err; @@ -4271,6 +4571,8 @@ export class PiAgent extends BaseAgent { // 身份注册(否则 ?session= ctx 泄漏)+ 关掉可能已 spawn 的子进程(否则僵尸 pi // 仍持有本会话的 MCP 路由),再把原始错误抛给调用方。 let sdkSessionId = ''; + const rpcReadyStartedAt = Date.now(); + let rpcReadyLogged = false; let runtimeCapabilityRefreshPromise: Promise | undefined; const refreshRuntimeCapabilities = async (stage: 'ready' | 'switch_session' | 'fork'): Promise => { const generation = ++runtimeCapabilityGeneration; @@ -4496,6 +4798,8 @@ export class PiAgent extends BaseAgent { } sdkSessionId = validateSdkSessionId(stateData.sessionFile || stateData.sessionId!); queue.push({ type: 'session_id', data: sdkSessionId, source: 'pi' }); + logPiStartupStage(this.deps.logger, startupTraceId, 'rpc-ready', rpcReadyStartedAt, 'ok'); + rpcReadyLogged = true; // get_state is the ready boundary. Capture exactly once after the final // fresh/resumed runtime has been selected; list/customization calls never @@ -4523,6 +4827,9 @@ export class PiAgent extends BaseAgent { } } } catch (err) { + if (!rpcReadyLogged) { + logPiStartupStage(this.deps.logger, startupTraceId, 'rpc-ready', rpcReadyStartedAt, 'degraded'); + } disposePiTranslateContext(ctx); try { disposeSessionRegistrations(); @@ -4538,19 +4845,20 @@ export class PiAgent extends BaseAgent { await proc.close(); } catch (error) { closeError = error; + if (!remote) configHomeCleanupDeferredToQuarantine = true; this.failedStartupCleanups.set(startupCleanupKey, { proc, promise: null, ...(!remote ? { cleanupLocal: () => { - cleanupConfigHome(); + void cleanupConfigHome(); cleanupRuntimeFiles(); } } : {}), }); } if (!remote && !closeError) { - cleanupConfigHome(); + await cleanupConfigHome(); cleanupRuntimeFiles(); } if (closeError) { @@ -4566,6 +4874,7 @@ export class PiAgent extends BaseAgent { this.launchSubagentRunner(request); const deps = this.deps; const agentKind = this.kind; + let firstModelRequestLogged = false; // 取消边界:main 的队列协调器在 Stop/close 抢占时会 abort 传入的 signal 并撤下 // steer 标记。send/steer 必须在**构建 prompt(读附件是 async)前后、投递 RPC 前** @@ -5162,6 +5471,9 @@ export class PiAgent extends BaseAgent { if (!managedPackageRoute.accepted) rejectIfCancelled(sendOpts, 'send'); const pendingTurnStartToken = markPiHostTurnStartPending(ctx); promptRequestStarted = true; + const firstModelRequestStartedAt = firstModelRequestLogged ? undefined : Date.now(); + if (firstModelRequestStartedAt !== undefined) firstModelRequestLogged = true; + let firstModelRequestStatus: 'ok' | 'degraded' = 'degraded'; try { doctorCommandActivity.enter(isDoctorCommand); const resp = await runExclusivePiRpc(() => proc.request(command, { @@ -5172,6 +5484,7 @@ export class PiAgent extends BaseAgent { refreshTimeoutOnEvent: (event) => PI_PROMPT_ACCEPTANCE_PROGRESS_EVENTS.has(event.type), })); + firstModelRequestStatus = resp.success ? 'ok' : 'degraded'; if (!resp.success) { rollbackPiHostTurnStart(ctx, pendingTurnStartToken); if (managedPackageRoute.accepted) { @@ -5297,6 +5610,15 @@ export class PiAgent extends BaseAgent { } } } finally { + if (firstModelRequestStartedAt !== undefined) { + logPiStartupStage( + deps.logger, + startupTraceId, + 'first-model-request', + firstModelRequestStartedAt, + firstModelRequestStatus, + ); + } if (activeExtensionCommandNotifications === capturedExtensionNotifications) { activeExtensionCommandNotifications = null; } @@ -5543,7 +5865,7 @@ export class PiAgent extends BaseAgent { piProcessExited = true; if (!remote) { // 会话结束:清理隔离的 configHome 与 runtime 文件(onExit 幂等,二者先到先清)。 - cleanupConfigHome(); + await cleanupConfigHome(); cleanupRuntimeFiles(); } }, @@ -5914,7 +6236,18 @@ export class PiAgent extends BaseAgent { }, }; + configHomeLifecycleTransferred = true; return handle; + } finally { + if ( + !remote + && localConfigHomeCreated + && !configHomeLifecycleTransferred + && !configHomeCleanupDeferredToQuarantine + ) { + await cleanupConfigHome(); + } + } } async getMemoryStatus(): Promise { diff --git a/packages/maker-core/src/agents/pi/pi-subagent-runs.ts b/packages/maker-core/src/agents/pi/pi-subagent-runs.ts index 5e3f88e94b..e6045a6a83 100644 --- a/packages/maker-core/src/agents/pi/pi-subagent-runs.ts +++ b/packages/maker-core/src/agents/pi/pi-subagent-runs.ts @@ -1027,6 +1027,11 @@ function ownProcessStartTimeSec(): number { return OWN_PROCESS_START_TIME_SEC; } +/** Process-incarnation stamp shared by other host-owned Pi runtime records. */ +export function piHostProcessStartTimeSec(): number { + return ownProcessStartTimeSec(); +} + export interface PiSubagentOwnerIdentity { pid: number; /** Absent on ids written before the start time was recorded. */ @@ -1224,7 +1229,7 @@ function probeProcessStartTimeSec(pid: number, now: number): number | null { * cannot detect reuse; only a fresh probe can, and a sweep-scoped memo is the * largest window in which reuse is not observable anyway. */ -type ProcessStartTimeMemo = Map; +export type ProcessStartTimeMemo = Map; function readProcessStartTimeSec(pid: number, memo?: ProcessStartTimeMemo): number | null { const cached = memo?.get(pid); @@ -1264,6 +1269,14 @@ function isOwnerInstanceAlive( return Math.abs(startTimeSec - identity.startTimeSec) <= OWNER_START_TIME_TOLERANCE_SEC; } +/** Conservative liveness check for a Pi host process incarnation. */ +export function isPiHostProcessInstanceAlive( + identity: PiSubagentOwnerIdentity, + startTimeMemo?: ProcessStartTimeMemo, +): boolean { + return isOwnerInstanceAlive(identity, startTimeMemo); +} + function isProcessAlive(pid: number | undefined): boolean | null { if (!Number.isSafeInteger(pid) || (pid ?? 0) <= 0) return null; try {