From 09841c8dbd7c58ce9ecda9dc49ea1c1b57759b67 Mon Sep 17 00:00:00 2001 From: brother7 <7brother7@gmail.com> Date: Mon, 24 Aug 2026 17:59:49 +0800 Subject: [PATCH] fix(pi): harden release proof lifecycle --- .../pi/release-proof-cleanup.ts | 38 +++++ electron/coding-runtime/pi/release-proof.ts | 142 ++++++++++++++---- electron/main/index.ts | 16 +- scripts/run-pi-release-performance.mjs | 1 + scripts/run-pi-subagent-packaged-smoke.mjs | 58 ++++++- scripts/smoke-pi-real.mjs | 1 + tests/unit/pi-release-proof-cleanup.test.ts | 66 ++++++++ 7 files changed, 286 insertions(+), 36 deletions(-) create mode 100644 electron/coding-runtime/pi/release-proof-cleanup.ts create mode 100644 tests/unit/pi-release-proof-cleanup.test.ts diff --git a/electron/coding-runtime/pi/release-proof-cleanup.ts b/electron/coding-runtime/pi/release-proof-cleanup.ts new file mode 100644 index 0000000..6d1f5ed --- /dev/null +++ b/electron/coding-runtime/pi/release-proof-cleanup.ts @@ -0,0 +1,38 @@ +export type PiReleasePressureCleanupStep = { + name: string; + run(): void | Promise; +}; + +export async function runPiReleasePressureCleanup( + steps: readonly PiReleasePressureCleanupStep[], + options: { stepTimeoutMs?: number } = {}, +): Promise { + const stepTimeoutMs = options.stepTimeoutMs ?? 10_000; + if (!Number.isSafeInteger(stepTimeoutMs) || stepTimeoutMs <= 0) { + throw new Error('PI release pressure cleanup timeout must be a positive safe integer'); + } + const failures: Array<{ name: string; error: unknown }> = []; + for (const step of steps) { + let timeout: ReturnType | undefined; + try { + await Promise.race([ + Promise.resolve().then(() => step.run()), + new Promise((_resolve, reject) => { + timeout = setTimeout(() => { + reject(new Error(`PI release pressure cleanup timed out: ${step.name}`)); + }, stepTimeoutMs); + }), + ]); + } catch (error) { + failures.push({ name: step.name, error }); + } finally { + if (timeout) clearTimeout(timeout); + } + } + if (failures.length > 0) { + throw new AggregateError( + failures.map(({ error }) => error), + `PI release pressure cleanup failed: ${failures.map(({ name }) => name).join(', ')}`, + ); + } +} diff --git a/electron/coding-runtime/pi/release-proof.ts b/electron/coding-runtime/pi/release-proof.ts index d0ef36f..7489a94 100644 --- a/electron/coding-runtime/pi/release-proof.ts +++ b/electron/coding-runtime/pi/release-proof.ts @@ -15,6 +15,7 @@ import type { ProviderAccount } from '../../shared/providers/types'; import type { PrepareConversationInput } from '../contracts'; import { PiManagedExtensionHost } from './extension-host'; import { PiManagedInputRevisionCoordinator } from './managed-input-revision'; +import { runPiReleasePressureCleanup } from './release-proof-cleanup'; import type { PiRpcEvent } from './rpc-client'; import { createPiManagedWorkerOpener } from './runtime'; import { PiSessionRegistry } from './session-registry'; @@ -28,6 +29,7 @@ import { PiProjectWriteLeaseCoordinator, type PiProjectWriteLease } from './writ type ProofWorkerRole = 'parent' | 'child'; type ProofProviderMode = 'subagent' | 'pressure'; type ProofMilestone = PiRuntimeTelemetryEvent['milestone'] | 'agent.start' | 'provider.first_event'; +type ProofMilestoneSource = 'main.telemetry' | 'pi.agent_start' | 'pi.assistant_message_start'; type TrackedProcess = { role: ProofWorkerRole; @@ -80,7 +82,7 @@ type RealProofComposition = { type PressureRun = { active: PiReleasePressureSnapshot; - finish(): Promise; + finish(options?: { injectFailureAt?: 'parents.settle' }): Promise; }; export interface PiReleaseManagedTurnProof { @@ -89,6 +91,7 @@ export interface PiReleaseManagedTurnProof { workerGeneration: number; milestones: Array<{ milestone: ProofMilestone; + source: ProofMilestoneSource; durationMs: number; at: number; }>; @@ -116,6 +119,7 @@ export interface PiReleaseExtensionProof { subagentStatus: string; subagentSummary: string; materializedExtension: string; + providerFirstEventDelayMs: number; managedTurns: PiReleaseManagedTurnProof[]; managedWorkerMilestones: PiRuntimeTelemetryEvent[]; released: { @@ -129,6 +133,7 @@ export interface PiReleaseExtensionProof { const PROOF_ACCOUNT_ID = 'release-proof-account'; const PROOF_AGENT_ID = 'release-proof-agent'; const PROOF_MODEL_ID = 'release-proof-model'; +const PROOF_PROVIDER_FIRST_EVENT_DELAY_MS = 75; const EXPECTED_TURN_MILESTONES: readonly ProofMilestone[] = [ 'worker.queue_wait', 'resources.ready', @@ -152,6 +157,19 @@ async function waitFor(predicate: () => boolean, message: string): Promise throw new Error(message); } +async function delay(durationMs: number): Promise { + await new Promise((resolve) => setTimeout(resolve, durationMs)); +} + +function isAssistantMessageStart(event: PiRpcEvent): boolean { + if (event.type !== 'message_start') return false; + const message = event.message; + return Boolean(message) + && typeof message === 'object' + && !Array.isArray(message) + && (message as { role?: unknown }).role === 'assistant'; +} + function shortRef(value: string): string { return value.replace(/[^A-Za-z0-9]/g, '').slice(-8).toLowerCase() || 'unknown'; } @@ -249,6 +267,8 @@ function respondWithSubagentCall(response: ServerResponse, model: string): void async function startLocalProofProvider(mode: ProofProviderMode): Promise { const requests: ProviderRequest[] = []; const held = new Set(); + let closed = false; + let closeFlight: Promise | null = null; const server: Server = createServer(async (request, response) => { response.once('error', () => undefined); try { @@ -264,10 +284,12 @@ async function startLocalProofProvider(mode: ProofProviderMode): Promise release('child'), releaseAll: () => release(), close: async () => { - release(); - await new Promise((resolve, reject) => { - server.close((error) => (error ? reject(error) : resolve())); - server.closeIdleConnections?.(); - }); + if (closed) return; + if (!closeFlight) { + closeFlight = (async () => { + release(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + server.closeIdleConnections?.(); + }); + closed = true; + })().finally(() => { + closeFlight = null; + }); + } + await closeFlight; }, }; } @@ -522,29 +553,35 @@ function timelineForTurn( const agentStart = observed.find(({ event }) => event.type === 'agent_start'); const agentSettled = [...observed].reverse().find(({ event }) => event.type === 'agent_settled'); const firstProviderEvent = agentStart && observed.find((entry) => ( - entry !== agentStart - && entry.at >= agentStart.at - && entry.event.type !== 'agent_settled' + entry.at >= agentStart.at && isAssistantMessageStart(entry.event) )); const promptAccepted = byMilestone.get('prompt.accepted'); if (!agentStart || !agentSettled || !firstProviderEvent || !promptAccepted) { throw new Error(`Managed ${cold ? 'cold' : 'warm'} turn did not expose provider lifecycle events`); } - const proofEvents = new Map(); + const proofEvents = new Map(); for (const event of runTelemetry) { proofEvents.set(event.milestone, { milestone: event.milestone, + source: 'main.telemetry', durationMs: event.durationMs, at: event.at, }); } proofEvents.set('agent.start', { milestone: 'agent.start', + source: 'pi.agent_start', durationMs: Math.max(0, agentStart.at - promptAccepted.at), at: agentStart.at, }); proofEvents.set('provider.first_event', { milestone: 'provider.first_event', + source: 'pi.assistant_message_start', durationMs: Math.max(0, firstProviderEvent.at - agentStart.at), at: firstProviderEvent.at, }); @@ -699,6 +736,7 @@ async function startPressureRun(): Promise { providerRequests: providerRequestCounts(provider), childDiagnostics, })}`, + { cause: error }, ); } const active = pressureSnapshot(composition, provider, writeLeases); @@ -709,22 +747,64 @@ async function startPressureRun(): Promise { } return { active, - finish: async () => { - provider.releaseAll(); - await Promise.all(dispatches); - await waitFor( - () => composition.pool.getDiagnostics().workers.every(({ state }) => state === 'idle'), - 'Persistent Pi parents did not settle after pressure release', - ); - await composition.scheduler.close(); - await composition.pool.shutdown(); - for (const lease of heldWriteLeases.splice(0)) lease.release(); - await composition.extensionHost.close(); - await provider.close(); - const released = pressureSnapshot(composition, provider, writeLeases); - await rm(root, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); - return released; - }, + finish: (() => { + let flight: Promise | null = null; + const cleanup = async (injectFailureAt?: 'parents.settle') => { + const steps = [ + { name: 'provider.release', run: () => provider.releaseAll() }, + { + name: 'dispatches.settle', + run: async () => { + const results = await Promise.allSettled(dispatches); + const failures = results.flatMap((result) => ( + result.status === 'rejected' ? [result.reason] : [] + )); + if (failures.length > 0) { + throw new AggregateError(failures, 'PI release pressure dispatch failed'); + } + }, + }, + { + name: 'parents.settle', + run: async () => await waitFor( + () => composition.pool.getDiagnostics().workers.every(({ state }) => state === 'idle'), + 'Persistent Pi parents did not settle after pressure release', + ), + }, + { name: 'scheduler.close', run: async () => await composition.scheduler.close() }, + { name: 'pool.shutdown', run: async () => await composition.pool.shutdown() }, + { + name: 'leases.release', + run: () => { + for (const lease of heldWriteLeases.splice(0)) lease.release(); + }, + }, + { name: 'extension.close', run: async () => await composition.extensionHost.close() }, + { name: 'provider.close', run: async () => await provider.close() }, + { + name: 'scratch.remove', + run: async () => await rm(root, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }), + }, + ]; + if (injectFailureAt) { + const target = steps.find(({ name }) => name === injectFailureAt); + if (!target) throw new Error(`Unknown PI release cleanup failure injection: ${injectFailureAt}`); + target.run = () => { + throw new Error(`Injected PI release cleanup failure at ${injectFailureAt}`); + }; + } + await runPiReleasePressureCleanup(steps); + return pressureSnapshot(composition, provider, writeLeases); + }; + return (options?: { injectFailureAt?: 'parents.settle' }) => { + if (!flight) { + flight = cleanup(options?.injectFailureAt).finally(() => { + flight = null; + }); + } + return flight; + }; + })(), }; } catch (error) { provider.releaseAll(); @@ -842,6 +922,7 @@ export async function runFinalAsarExtensionProof(): Promise { +export async function finishFinalAsarPressureProof( + options?: { injectFailureAt?: 'parents.settle' }, +): Promise { const current = pressureRun; if (!current) throw new Error('PI release pressure proof is not running'); - pressureRun = null; - return await current.finish(); + const released = await current.finish(options); + if (pressureRun === current) pressureRun = null; + return released; } diff --git a/electron/main/index.ts b/electron/main/index.ts index ffde4a0..b577247 100644 --- a/electron/main/index.ts +++ b/electron/main/index.ts @@ -880,7 +880,11 @@ export async function runLocalPreviewPreflightE2E( } } -type PiReleaseProofAction = 'extension' | 'pressure.start' | 'pressure.finish'; +type PiReleaseProofAction = + | 'extension' + | 'pressure.start' + | 'pressure.finish' + | 'pressure.finish.inject-failure'; export async function runPiReleaseProofE2E(action: PiReleaseProofAction) { if (!isE2EMode) throw new Error('PI release proof is unavailable'); @@ -896,7 +900,15 @@ export async function runPiReleaseProofE2E(action: PiReleaseProofAction) { if (action === 'pressure.start') { return { action, packagedMain, pressure: await startFinalAsarPressureProof() }; } - return { action, packagedMain, pressure: await finishFinalAsarPressureProof() }; + return { + action, + packagedMain, + pressure: await finishFinalAsarPressureProof( + action === 'pressure.finish.inject-failure' + ? { injectFailureAt: 'parents.settle' } + : undefined, + ), + }; } if (isE2EMode) { diff --git a/scripts/run-pi-release-performance.mjs b/scripts/run-pi-release-performance.mjs index 3bc91e7..2de6096 100644 --- a/scripts/run-pi-release-performance.mjs +++ b/scripts/run-pi-release-performance.mjs @@ -133,6 +133,7 @@ async function runProductProofSamples(artifact, projectRoot, samples) { runtimeRoot: artifact.artifact.runtimeRoot, electronExecutable: artifact.artifact.executable, reportPath: undefined, + verifyCleanupFailure: false, })); } return reports; diff --git a/scripts/run-pi-subagent-packaged-smoke.mjs b/scripts/run-pi-subagent-packaged-smoke.mjs index 32caf9a..b1c13b3 100644 --- a/scripts/run-pi-subagent-packaged-smoke.mjs +++ b/scripts/run-pi-subagent-packaged-smoke.mjs @@ -16,6 +16,7 @@ function parseArgs(argv, projectRoot = process.cwd()) { runtimeRoot: undefined, electronExecutable: undefined, reportPath: undefined, + verifyCleanupFailure: true, }; for (let index = 0; index < argv.length; index += 1) { const argument = argv[index]; @@ -107,6 +108,27 @@ function assertReleasedPressure(pressure) { } } +function assertManagedTurns(extension) { + const expectedMilestones = 'worker.queue_wait,resources.ready,worker.spawn,rpc.ready,session.open,prompt.accepted,agent.start,provider.first_event,agent.settled'; + const expectedDelayMs = extension?.providerFirstEventDelayMs; + if (!Number.isFinite(expectedDelayMs) || expectedDelayMs < 1) { + throw new Error(`PI proof did not expose a controlled Provider first-event delay: ${expectedDelayMs}`); + } + for (const turn of extension?.managedTurns ?? []) { + if (turn.milestones.map(({ milestone }) => milestone).join(',') !== expectedMilestones) { + throw new Error(`PI managed turn timeline is incomplete: ${JSON.stringify(turn)}`); + } + const agentStart = turn.milestones.find(({ milestone }) => milestone === 'agent.start'); + const providerFirstEvent = turn.milestones.find(({ milestone }) => milestone === 'provider.first_event'); + if (agentStart?.source !== 'pi.agent_start' + || providerFirstEvent?.source !== 'pi.assistant_message_start' + || providerFirstEvent.at <= agentStart.at + || providerFirstEvent.durationMs < Math.floor(expectedDelayMs * 0.75)) { + throw new Error(`PI Provider first-event milestone is not Provider-response-backed: ${JSON.stringify(turn)}`); + } + } +} + async function evaluateProof(electronApplication, action) { return await electronApplication.evaluate(async (_electron, requestedAction) => { const proof = globalThis.__niancodeRunPiReleaseProofE2E; @@ -158,17 +180,14 @@ export async function runPackagedProductProof(options) { const extension = await evaluateProof(electronApplication, 'extension'); assertPackagedMain(extension); + assertManagedTurns(extension.extension); if (extension.extension?.subagentStatus !== 'complete' || extension.extension?.childToolNames?.join(',') !== 'find,grep,ls,read' || !extension.extension?.parentToolNames?.includes('subagent') || extension.extension?.parentProcessIds?.length < 2 || extension.extension?.childProcessIds?.length !== 1 || extension.extension?.providerRequests?.child !== 1 - || extension.extension?.managedTurns?.length !== 2 - || extension.extension.managedTurns.some(({ milestones }) => ( - milestones.map(({ milestone }) => milestone).join(',') - !== 'worker.queue_wait,resources.ready,worker.spawn,rpc.ready,session.open,prompt.accepted,agent.start,provider.first_event,agent.settled' - ))) { + || extension.extension?.managedTurns?.length !== 2) { throw new Error(`Final ASAR extension/subagent proof failed: ${JSON.stringify(extension.extension)}`); } @@ -187,6 +206,34 @@ export async function runPackagedProductProof(options) { assertPackagedMain(pressureFinish); assertReleasedPressure(pressureFinish.pressure); + let failureCleanup; + if (options.verifyCleanupFailure === true) { + const failureStart = await evaluateProof(electronApplication, 'pressure.start'); + pressureActive = true; + assertPackagedMain(failureStart); + assertActivePressure(failureStart.pressure); + let injectedFailure; + try { + await evaluateProof(electronApplication, 'pressure.finish.inject-failure'); + } catch (error) { + injectedFailure = error; + } + if (!injectedFailure + || !String(injectedFailure).includes('PI release pressure cleanup failed: parents.settle')) { + throw new Error(`PI pressure cleanup failure injection did not fail as expected: ${String(injectedFailure)}`); + } + const retryFinish = await evaluateProof(electronApplication, 'pressure.finish'); + pressureActive = false; + assertPackagedMain(retryFinish); + assertReleasedPressure(retryFinish.pressure); + failureCleanup = { + injectedAt: 'parents.settle', + firstFinish: 'failed-as-injected', + retry: 'pass', + released: retryFinish.pressure, + }; + } + const report = { schemaVersion: 1, generatedAt: new Date().toISOString(), @@ -203,6 +250,7 @@ export async function runPackagedProductProof(options) { durationMs: uiInteractiveMs, }, released: pressureFinish.pressure, + ...(failureCleanup ? { failureCleanup } : {}), }, result: 'pass', }; diff --git a/scripts/smoke-pi-real.mjs b/scripts/smoke-pi-real.mjs index a8b2d60..4d1bcca 100644 --- a/scripts/smoke-pi-real.mjs +++ b/scripts/smoke-pi-real.mjs @@ -35,6 +35,7 @@ async function runExtensionSmoke(projectRoot, artifact) { runtimeRoot: artifact.artifact.runtimeRoot, electronExecutable: artifact.artifact.executable, reportPath: undefined, + verifyCleanupFailure: true, }); } diff --git a/tests/unit/pi-release-proof-cleanup.test.ts b/tests/unit/pi-release-proof-cleanup.test.ts new file mode 100644 index 0000000..de30cd8 --- /dev/null +++ b/tests/unit/pi-release-proof-cleanup.test.ts @@ -0,0 +1,66 @@ +import { describe, expect, it } from 'vitest'; + +import { runPiReleasePressureCleanup } from '../../electron/coding-runtime/pi/release-proof-cleanup'; + +describe('PI release pressure cleanup', () => { + it('continues after an injected failure, reaches zero, and can be retried', async () => { + const counts = { + parentWorkers: 4, + childWorkers: 4, + liveProcessIds: 8, + providerRequests: 8, + processBudget: 8, + childPermits: 4, + dispatches: 4, + writeLeases: 4, + }; + let failParentSettle = true; + const steps = [ + { name: 'provider.release', run: () => { counts.providerRequests = 0; } }, + { + name: 'dispatches.settle', + run: () => { + counts.childWorkers = 0; + counts.childPermits = 0; + counts.dispatches = 0; + }, + }, + { + name: 'parents.settle', + run: () => { + if (failParentSettle) { + failParentSettle = false; + throw new Error('injected parent settle failure'); + } + }, + }, + { + name: 'pool.shutdown', + run: () => { + counts.parentWorkers = 0; + counts.liveProcessIds = 0; + counts.processBudget = 0; + }, + }, + { name: 'leases.release', run: () => { counts.writeLeases = 0; } }, + ]; + + await expect(runPiReleasePressureCleanup(steps)) + .rejects.toThrow('PI release pressure cleanup failed: parents.settle'); + expect(Object.values(counts)).toEqual([0, 0, 0, 0, 0, 0, 0, 0]); + + await expect(runPiReleasePressureCleanup(steps)).resolves.toBeUndefined(); + expect(Object.values(counts)).toEqual([0, 0, 0, 0, 0, 0, 0, 0]); + }); + + it('bounds a hanging step and still runs later cleanup', async () => { + let released = false; + await expect(runPiReleasePressureCleanup([ + { name: 'hang', run: async () => await new Promise(() => undefined) }, + { name: 'release', run: () => { released = true; } }, + ], { stepTimeoutMs: 10 })).rejects.toThrow( + 'PI release pressure cleanup failed: hang', + ); + expect(released).toBe(true); + }); +});