fix(pi): harden release proof lifecycle

This commit is contained in:
2026-08-24 17:59:49 +08:00
parent ead9d1dbaf
commit 09841c8dbd
7 changed files with 286 additions and 36 deletions

View File

@@ -0,0 +1,38 @@
export type PiReleasePressureCleanupStep = {
name: string;
run(): void | Promise<void>;
};
export async function runPiReleasePressureCleanup(
steps: readonly PiReleasePressureCleanupStep[],
options: { stepTimeoutMs?: number } = {},
): Promise<void> {
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<typeof setTimeout> | undefined;
try {
await Promise.race([
Promise.resolve().then(() => step.run()),
new Promise<never>((_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(', ')}`,
);
}
}

View File

@@ -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<PiReleasePressureSnapshot>;
finish(options?: { injectFailureAt?: 'parents.settle' }): Promise<PiReleasePressureSnapshot>;
};
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<void>
throw new Error(message);
}
async function delay(durationMs: number): Promise<void> {
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<LocalProofProvider> {
const requests: ProviderRequest[] = [];
const held = new Set<HeldProviderResponse>();
let closed = false;
let closeFlight: Promise<void> | null = null;
const server: Server = createServer(async (request, response) => {
response.once('error', () => undefined);
try {
@@ -264,10 +284,12 @@ async function startLocalProofProvider(mode: ProofProviderMode): Promise<LocalPr
requests.push({ role, toolNames, hasToolResult: toolResult });
if (mode === 'subagent' && role === 'parent' && !toolResult) {
await delay(PROOF_PROVIDER_FIRST_EVENT_DELAY_MS);
respondWithSubagentCall(response, model);
return;
}
if (mode === 'subagent' && role === 'parent') {
await delay(PROOF_PROVIDER_FIRST_EVENT_DELAY_MS);
respondWithText(response, model, 'REAL_PARENT_COMPLETE');
return;
}
@@ -314,11 +336,20 @@ async function startLocalProofProvider(mode: ProofProviderMode): Promise<LocalPr
releaseChildren: () => release('child'),
releaseAll: () => release(),
close: async () => {
release();
await new Promise<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
server.closeIdleConnections?.();
});
if (closed) return;
if (!closeFlight) {
closeFlight = (async () => {
release();
await new Promise<void>((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<ProofMilestone, { milestone: ProofMilestone; durationMs: number; at: number }>();
const proofEvents = new Map<ProofMilestone, {
milestone: ProofMilestone;
source: ProofMilestoneSource;
durationMs: number;
at: number;
}>();
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<PressureRun> {
providerRequests: providerRequestCounts(provider),
childDiagnostics,
})}`,
{ cause: error },
);
}
const active = pressureSnapshot(composition, provider, writeLeases);
@@ -709,22 +747,64 @@ async function startPressureRun(): Promise<PressureRun> {
}
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<PiReleasePressureSnapshot> | 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<PiReleaseExtensionPr
subagentStatus: 'complete',
subagentSummary: 'REAL_CHILD_COMPLETE',
materializedExtension: 'makelore-runtime-v3.mjs',
providerFirstEventDelayMs: PROOF_PROVIDER_FIRST_EVENT_DELAY_MS,
managedTurns,
managedWorkerMilestones: composition.telemetry,
released,
@@ -854,9 +935,12 @@ export async function startFinalAsarPressureProof(): Promise<PiReleasePressureSn
return pressureRun.active;
}
export async function finishFinalAsarPressureProof(): Promise<PiReleasePressureSnapshot> {
export async function finishFinalAsarPressureProof(
options?: { injectFailureAt?: 'parents.settle' },
): Promise<PiReleasePressureSnapshot> {
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;
}

View File

@@ -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) {