Files
makelore/electron/coding-runtime/pi/runtime.ts
brother7 0f7093d173 fix: allow replacing unavailable conversation models
Persist the selected model before opening dormant sessions, and reconfigure crashed workers on the same account. Preserve session history and report model removal through the model-unavailable contract instead of a generic runtime failure.
2026-09-21 12:23:04 +08:00

2028 lines
75 KiB
TypeScript

import { randomUUID } from 'node:crypto';
import path from 'node:path';
import { logger } from '../../utils/logger';
import type { ModelSummary, ProviderAccount } from '../../shared/providers/types';
import { validateSessionKey } from '../../coding-projects/conversation-store';
import {
buildPiProviderCatalog,
buildPiWorkerCredentialProjection,
selectPiProviderModel,
writePiProviderCatalog,
type PiProviderSelection,
} from './provider-config';
import { PiProviderRefreshCoordinator } from './provider-refresh';
import type {
CodingConversationRuntime,
CodingRuntimeDisposeReason,
CodingRuntimeCommand,
CodingRuntimeDiagnostics,
CodingRuntimePublicError,
ConversationInteraction,
ConversationInteractionResponse,
ConversationModelState,
ConversationPatch,
ConversationPatchEnvelope,
ConversationRuntimeState,
ConversationSnapshot,
ForkConversationInput,
ForkResult,
PrepareConversationInput,
ProductModelRef,
PromptAcceptance,
PromptConversationInput,
QueueAcceptance,
QueueMessageInput,
SetConversationModelInput,
SetThinkingLevelInput,
} from '../contracts';
import {
createConversationReducerState,
reduceConversationPatch,
type ConversationReducerState,
} from '../conversation-reducer';
import { CodingRuntimeContractError } from '../runtime-errors';
import {
PiEventProjector,
type PiEventProjectorOptions,
} from './event-projector';
import { PiProcessError } from './process-errors';
import type { PiSessionRegistry } from './session-registry';
import {
buildPiManagedInputArgs,
ensurePiManagedPaths,
materializePiAgentResources,
} from './resource-loader';
import type {
PiRpcCommand,
PiRpcEvent,
PiRpcRequestOptions,
PiRpcResponse,
} from './rpc-client';
import {
PiWorkerProcess,
PI_CORE_TOOL_NAMES,
type PiWorkerProcessOptions,
type PiWorkerProofFailure,
type PiWorkerStopReason,
type PiWorkerStopResult,
} from './worker-process';
import {
createPiRuntimeTelemetryEvent,
type PiRuntimeMilestone,
type PiRuntimeTelemetryEvent,
} from './telemetry';
import {
type PiConversationWorker,
type PiWorkerOpenInput,
type PiWorkerOpenResult,
type PiWorkerPoolEvent,
type PiWorkerPoolState,
PiWorkerPool,
} from './worker-pool';
import {
PiSessionProjectionError,
projectPiSessionSnapshot,
} from './session-projector';
import { PiManagedExtensionHost } from './extension-host';
import type { PiSubagentScheduler } from './subagent';
import type { ModelToolRegistryPort } from './model-tools/model-tool-registry';
import type { DevicePackageManager } from '../../coding-packages/device-package-manager';
import type { CodingPluginToolDefinition } from '../../../shared/coding-plugins';
import type { CodingCapabilityRegistry, ResolvedWorkerResources } from '../../coding-plugins/registry';
import {
PiInteractionStore,
} from './interaction';
import {
PiExtensionUiProjector,
type PiExtensionUiProjection,
} from './extension-ui-projector';
type RuntimeIdKind = 'run' | 'queue';
import { buildManagedModelRequest, managedPiThinkingLevel, validateManagedReasoningChoice } from '../../../shared/managed-model-capabilities';
export interface PiConversationRuntimeOptions {
pool: PiWorkerPool;
registry: PiSessionRegistry;
resolveModel(model: ProductModelRef): Promise<PiProviderSelection>;
resolveImages?(attachments: Array<{ attachmentId: string }>): Promise<Array<{
type: 'image'; data: string; mimeType: string;
}>>;
projectImage?: PiEventProjectorOptions['projectImage'];
createId?(kind: RuntimeIdKind): string;
now?: () => number;
providerRefreshCoordinator?: PiProviderRefreshCoordinator;
isAuthenticationError?(error: unknown): boolean;
refreshCredential?(accountId: string): Promise<void>;
extensionHost?: PiManagedExtensionHost;
subagentScheduler?: PiSubagentScheduler;
getDraftRevision?(conversationId: string): number;
knownExtensionWidgetKeys?: readonly string[];
onExtensionUiProjection?(projection: PiExtensionUiProjection): void;
acquireBackgroundLease?(lease: { id: string; kind: 'coding-run' }): () => void;
}
export interface PiWorkerProcessAdapter {
readonly generation: number;
start(): Promise<unknown>;
request<T = unknown>(
command: PiRpcCommand,
options?: PiRpcRequestOptions,
): Promise<PiRpcResponse<T>>;
send(command: PiRpcCommand): Promise<void>;
subscribe(listener: (event: PiRpcEvent) => void): () => void;
subscribeInvalidation(listener: (error: PiProcessError) => void): () => void;
stop(reason: PiWorkerStopReason): Promise<PiWorkerStopResult>;
injectFailureForProof?(failure: PiWorkerProofFailure): Promise<void>;
delayNextResponseForProof?(commandType: string, delayMs: number): void;
}
export interface PiManagedProviderInput {
accounts: ProviderAccount[];
modelSummaries: ModelSummary[];
}
export interface PiManagedWorkerOpenerOptions {
registry: PiSessionRegistry;
executablePath: string;
cliPath: string;
userDataDir: string;
bundledSkillsDir: string;
getSkillRoots?(): readonly string[] | Promise<readonly string[]>;
registerActivePluginReleases?(releaseIds: readonly string[]): () => void | Promise<void>;
loadProviderInput(): Promise<PiManagedProviderInput>;
resolveCredential(account: ProviderAccount): Promise<string | null>;
getLocalProxyCredential?(): Promise<string | undefined>;
createSessionKey?: () => string;
createProcess?: (options: PiWorkerProcessOptions) => PiWorkerProcessAdapter;
now?: () => number;
onTelemetry?: (event: PiRuntimeTelemetryEvent) => void;
extensionHost: PiManagedExtensionHost;
capabilityRegistry?: CodingCapabilityRegistry;
modelToolRegistry?: ModelToolRegistryPort;
devicePackageManager?: Pick<DevicePackageManager, 'resolveEnabledResources' | 'registerActiveWorker'>;
devicePackageTools?: readonly CodingPluginToolDefinition[];
}
interface PiRpcSessionStateProjection {
sessionId: string;
sessionFile?: string;
}
function fallbackWorkerResources(_skillIds: readonly string[]): ResolvedWorkerResources {
return {
catalogRevision: 0,
pluginIds: [],
effectiveSkillIds: [],
skillEntries: [],
tools: [],
};
}
class ManagedPiConversationWorker implements PiConversationWorker {
constructor(
readonly id: string,
readonly generation: number,
private readonly process: PiWorkerProcessAdapter,
private readonly disposeExtension: () => Promise<void>,
private readonly unsubscribeExtensionInvalidation: () => void,
) {}
request<T = unknown>(command: PiRpcCommand, options?: PiRpcRequestOptions): Promise<PiRpcResponse<T>> {
return this.process.request<T>(command, options);
}
send(command: PiRpcCommand): Promise<void> {
return this.process.send(command);
}
subscribe(listener: (event: PiRpcEvent) => void): () => void {
return this.process.subscribe(listener);
}
subscribeInvalidation(listener: (error: PiProcessError) => void): () => void {
return this.process.subscribeInvalidation(listener);
}
async stop(reason: PiWorkerStopReason): Promise<PiWorkerStopResult> {
this.unsubscribeExtensionInvalidation();
try {
return await this.process.stop(reason);
} finally {
await this.disposeExtension();
}
}
async injectFailureForProof(failure: PiWorkerProofFailure): Promise<void> {
if (!this.process.injectFailureForProof) {
throw new Error('Managed Pi process does not support proof failure injection');
}
await this.process.injectFailureForProof(failure);
}
delayNextResponseForProof(commandType: string, delayMs: number): void {
if (!this.process.delayNextResponseForProof) {
throw new Error('Managed Pi process does not support proof response delay');
}
this.process.delayNextResponseForProof(commandType, delayMs);
}
}
export function createPiManagedWorkerOpener(
options: PiManagedWorkerOpenerOptions,
): (input: PiWorkerOpenInput) => Promise<PiWorkerOpenResult> {
const createProcess = options.createProcess ?? ((processOptions) => new PiWorkerProcess(processOptions));
const createSessionKey = options.createSessionKey ?? randomUUID;
const now = options.now ?? Date.now;
return async (input) => {
const resourcesStartedAt = now();
const registered = await options.registry.prepare(input.conversation);
const existingSession = input.existingSession ?? registered.session ?? undefined;
const model = registered.conversation.model;
if (!model || registered.conversation.modelResolution !== 'resolved') {
throw new CodingRuntimeContractError(
'CODING_MIGRATION_MODEL_REQUIRED',
'Conversation model must be selected before opening Pi',
true,
);
}
const providerInput = await options.loadProviderInput();
const catalog = buildPiProviderCatalog(providerInput);
const selection = selectPiProviderModel(catalog, model);
const account = providerInput.accounts.find((candidate) => candidate.id === selection.accountId);
const descriptor = catalog.descriptors.find((candidate) => candidate.accountId === selection.accountId);
if (!account || !descriptor) throw new Error('Selected Provider account is unavailable');
const managedPaths = await ensurePiManagedPaths(options.userDataDir);
await writePiProviderCatalog(managedPaths.modelsFile, catalog, model);
const workerResources = options.capabilityRegistry
? await options.capabilityRegistry.resolveWorkerResources({
projectPath: registered.projectPath,
projectId: input.conversation.projectId,
assignedSkillIds: registered.agent.skillIds,
role: 'parent',
})
: fallbackWorkerResources(registered.agent.skillIds);
const deviceResources = options.devicePackageManager
? await options.devicePackageManager.resolveEnabledResources()
: undefined;
const combinedSkillEntries = [
...workerResources.skillEntries,
...(deviceResources?.skillEntries ?? []),
];
const resources = await materializePiAgentResources({
userDataDir: options.userDataDir,
projectId: input.conversation.projectId,
agentId: registered.agent.id,
prompt: registered.agent.prompt,
skillEntries: combinedSkillEntries,
catalogRevision: workerResources.catalogRevision,
bundledSkillsDir: options.bundledSkillsDir,
...(workerResources.skillRoots
? { skillRoots: workerResources.skillRoots }
: options.getSkillRoots ? { skillRoots: await options.getSkillRoots() } : {}),
...(workerResources.effectiveSnapshot
? { effectiveSnapshot: workerResources.effectiveSnapshot }
: {}),
...(deviceResources
? {
devicePackageGeneration: deviceResources.generation,
devicePackageIds: deviceResources.packageIds,
}
: {}),
revision: input.revision,
});
const credential = await buildPiWorkerCredentialProjection({
account,
descriptor,
resolveCredential: options.resolveCredential,
...(options.getLocalProxyCredential
? { localProxyCredential: await options.getLocalProxyCredential() }
: {}),
});
const modelTools = options.modelToolRegistry?.registerWorker({
conversationId: input.conversation.conversationId,
generation: input.generation,
account,
descriptor,
selection,
credential,
});
let extension;
try {
extension = await options.extensionHost.registerWorker({
...(selection.managedCapability ? { managedModelRequest: buildManagedModelRequest(
model.modelId, model.reasoningChoice ?? { mode: 'default' }, selection.managedCapability,
) } : {}),
conversationId: input.conversation.conversationId,
generation: input.generation,
projectId: input.conversation.projectId,
projectPath: registered.projectPath,
skillEntries: combinedSkillEntries,
catalogRevision: workerResources.catalogRevision,
tools: [
...workerResources.tools,
...(modelTools?.tools ?? []),
...(options.devicePackageTools ?? []),
],
...(workerResources.effectiveSnapshot
? { effectiveSnapshot: workerResources.effectiveSnapshot }
: {}),
...(deviceResources
? {
devicePackageGeneration: deviceResources.generation,
devicePackageIds: deviceResources.packageIds,
}
: {}),
extensionsDir: managedPaths.extensionsDir,
});
} catch (error) {
modelTools?.dispose();
throw error;
}
recordManagedMilestone(
options.onTelemetry,
input,
'resources.ready',
now() - resourcesStartedAt,
now(),
!existingSession,
);
const sessionKey = validateSessionKey(
existingSession?.sessionKey ?? createSessionKey(),
);
if (existingSession && input.fork) {
throw new Error('Pi worker cannot reopen and fork a session at the same time');
}
const process = createProcess({
executablePath: options.executablePath,
cliPath: options.cliPath,
cwd: registered.projectPath,
configDir: resources.paths.configDir,
sessionDir: resources.projectSessionsDir,
tools: [...new Set([
...PI_CORE_TOOL_NAMES,
...extension.allowedToolNames,
])],
additionalArgs: [
...buildPiManagedInputArgs(selection, resources),
'--extension', extension.extensionPath,
...(deviceResources?.extensionPaths.flatMap((extensionPath) => [
'--extension', extensionPath,
]) ?? []),
...(input.fork ? ['--fork', input.fork.sourceSession.piSessionId] : []),
'--session-id', sessionKey,
],
env: { ...credential.env, ...extension.env },
sensitiveValues: [...credential.sensitiveValues, ...extension.sensitiveValues],
conversationId: input.conversation.conversationId,
workerGeneration: input.generation,
});
const releaseActivePluginReleases = workerResources.effectiveSnapshot
? options.registerActivePluginReleases?.(workerResources.effectiveSnapshot.pluginReleaseIds)
: undefined;
const releaseActiveDevicePackages = deviceResources
? options.devicePackageManager?.registerActiveWorker(deviceResources.packageRefs)
: undefined;
let managedResourcesDisposed = false;
const disposeManagedResources = async (): Promise<void> => {
if (managedResourcesDisposed) return;
managedResourcesDisposed = true;
try {
await extension.dispose();
} finally {
try {
modelTools?.dispose();
} finally {
try {
await releaseActiveDevicePackages?.();
} finally {
await releaseActivePluginReleases?.();
}
}
}
};
let unsubscribeExtensionInvalidation = process.subscribeInvalidation(() => {
void disposeManagedResources();
});
try {
const spawnStartedAt = now();
await process.start();
recordManagedMilestone(
options.onTelemetry,
input,
'worker.spawn',
now() - spawnStartedAt,
now(),
!existingSession,
);
if (input.fork?.sourceEntryId) {
await process.request({ type: 'fork', entryId: input.fork.sourceEntryId });
}
const readyStartedAt = now();
const response = await process.request<PiRpcSessionStateProjection>(
{ type: 'get_state' },
{ retry: 'read-only-once' },
);
recordManagedMilestone(
options.onTelemetry,
input,
'rpc.ready',
now() - readyStartedAt,
now(),
!existingSession,
);
const piSessionId = response.data?.sessionId?.trim();
if (!piSessionId) throw new Error('Pi worker did not return a session id');
if (existingSession && existingSession.piSessionId !== piSessionId) {
throw new CodingRuntimeContractError(
'CODING_SESSION_UNREADABLE',
'Pi reopened a different Conversation session',
true,
);
}
if (response.data?.sessionFile) {
const sessionsRoot = path.resolve(resources.projectSessionsDir);
const sessionFile = path.resolve(response.data.sessionFile);
const relativeSessionFile = path.relative(sessionsRoot, sessionFile);
if (!relativeSessionFile || relativeSessionFile.startsWith('..') || path.isAbsolute(relativeSessionFile)) {
throw new CodingRuntimeContractError(
'CODING_SESSION_UNREADABLE',
'Pi session file is outside the managed session directory',
true,
);
}
}
const sessionStartedAt = now();
const bound = await options.registry.ensureBinding(input.conversation, async () => ({
piSessionId,
sessionKey,
}));
if (!bound.session || bound.session.piSessionId !== piSessionId || bound.session.sessionKey !== sessionKey) {
throw new CodingRuntimeContractError(
'CODING_SESSION_UNREADABLE',
'Pi session binding does not match the Conversation registry',
true,
);
}
recordManagedMilestone(
options.onTelemetry,
input,
'session.open',
now() - sessionStartedAt,
now(),
!existingSession,
);
return {
worker: new ManagedPiConversationWorker(
`${input.conversation.conversationId}:${input.generation}`,
input.generation,
process,
disposeManagedResources,
() => {
unsubscribeExtensionInvalidation();
unsubscribeExtensionInvalidation = () => undefined;
},
),
session: clone(bound.session),
};
} catch (error) {
unsubscribeExtensionInvalidation();
await process.stop('open_failure').catch(() => undefined);
await disposeManagedResources();
throw error;
}
};
}
function recordManagedMilestone(
listener: ((event: PiRuntimeTelemetryEvent) => void) | undefined,
input: PiWorkerOpenInput,
milestone: Extract<PiRuntimeMilestone, 'resources.ready' | 'worker.spawn' | 'rpc.ready' | 'session.open'>,
durationMs: number,
at: number,
cold: boolean,
): void {
listener?.(createPiRuntimeTelemetryEvent({
milestone,
conversationId: input.conversation.conversationId,
workerGeneration: input.generation,
cold,
durationMs,
at,
}));
}
function clone<T>(value: T): T {
return structuredClone(value);
}
const PRODUCT_THINKING_LEVELS = new Set<ProductModelRef['thinkingLevel']>([
'off',
'minimal',
'low',
'medium',
'high',
'max',
]);
function productThinkingLevel(value: unknown): ProductModelRef['thinkingLevel'] | null {
return typeof value === 'string'
&& PRODUCT_THINKING_LEVELS.has(value as ProductModelRef['thinkingLevel'])
? value as ProductModelRef['thinkingLevel']
: null;
}
function availableThinkingLevels(value: unknown): ProductModelRef['thinkingLevel'][] {
if (!value || typeof value !== 'object' || Array.isArray(value)) return [];
const levels = (value as { levels?: unknown }).levels;
if (!Array.isArray(levels)) return [];
return [...new Set(levels.flatMap((level) => {
const normalized = productThinkingLevel(level);
return normalized ? [normalized] : [];
}))];
}
function effectiveThinkingLevel(value: unknown): ProductModelRef['thinkingLevel'] | null {
if (!value || typeof value !== 'object' || Array.isArray(value)) return null;
return productThinkingLevel((value as { thinkingLevel?: unknown }).thinkingLevel);
}
function runIsTerminal(status: ConversationSnapshot['run']['status']): boolean {
return status === 'idle' || status === 'error';
}
const REQUEST_UNCERTAIN_MESSAGE = '请求确认延迟,可能仍在执行。请等待结果,或中止/恢复后再重试。';
function requestUncertainError(): CodingRuntimeContractError {
return new CodingRuntimeContractError(
'CODING_REQUEST_UNCERTAIN',
REQUEST_UNCERTAIN_MESSAGE,
true,
);
}
function publicWorkerState(state: PiWorkerPoolState): ConversationSnapshot['worker'] {
if (state.state === 'spawning') return { status: 'starting', generation: state.generation };
if (state.state === 'crashed') {
return {
status: 'error',
generation: state.generation,
error: {
code: state.failureCode === 'PI_RPC_PROTOCOL_ERROR'
? 'CODING_RUNTIME_PROTOCOL_ERROR'
: 'CODING_RUNTIME_START_FAILED',
message: '本地 Agent 服务已中断。',
recoverable: true,
},
};
}
return { status: 'ready', generation: state.generation };
}
function samePublicWorkerState(
left: ConversationSnapshot['worker'],
right: ConversationSnapshot['worker'],
): boolean {
return left.status === right.status
&& left.generation === right.generation
&& left.error?.code === right.error?.code
&& left.error?.message === right.error?.message
&& left.error?.recoverable === right.error?.recoverable;
}
function runtimeFailure(error: unknown): CodingRuntimePublicError {
if (error instanceof CodingRuntimeContractError) return clone(error.publicError);
if (error instanceof PiSessionProjectionError) {
return { code: error.code, message: error.message, recoverable: error.recoverable };
}
if (error instanceof PiProcessError) {
if (error.code === 'PI_RPC_PROTOCOL_ERROR') {
return {
code: 'CODING_RUNTIME_PROTOCOL_ERROR',
message: '本地 Agent 通信异常,当前对话已停止。',
recoverable: true,
};
}
if (error.code === 'PI_RPC_TIMEOUT') {
return {
code: 'CODING_REQUEST_UNCERTAIN',
message: REQUEST_UNCERTAIN_MESSAGE,
recoverable: true,
};
}
if (error.code === 'PI_RPC_EXITED') {
return {
code: 'CODING_RUNTIME_START_FAILED',
message: '本地 Agent 已中断,原请求未自动重发。',
recoverable: true,
};
}
}
return {
code: 'CODING_RUNTIME_START_FAILED',
message: '本地 Agent 服务暂不可用。',
recoverable: true,
};
}
function emptySnapshot(
input: PrepareConversationInput,
workerState: PiWorkerPoolState,
): ConversationSnapshot {
return {
schemaVersion: 1,
conversation: {
id: input.conversationId,
projectId: input.projectId,
agentId: input.agentId,
title: input.title,
model: clone(input.model),
},
nodes: [],
run: { status: 'idle' },
queue: { items: [] },
context: { usedTokens: 0, contextWindow: 0, compaction: 'idle' },
pendingInteractions: [],
worker: publicWorkerState(workerState),
cursor: { workerGeneration: workerState.generation, seq: 0 },
};
}
export class PiConversationRuntime implements CodingConversationRuntime {
private readonly pool: PiWorkerPool;
private readonly registry: PiSessionRegistry;
private readonly resolveModel: PiConversationRuntimeOptions['resolveModel'];
private readonly resolveImages: NonNullable<PiConversationRuntimeOptions['resolveImages']>;
private readonly projectImage: PiEventProjectorOptions['projectImage'];
private readonly createRuntimeId: (kind: RuntimeIdKind) => string;
private readonly now: () => number;
private readonly providerRefresh: PiProviderRefreshCoordinator;
private readonly isAuthenticationError: ((error: unknown) => boolean) | undefined;
private readonly refreshCredential: ((accountId: string) => Promise<void>) | undefined;
private readonly extensionHost: PiManagedExtensionHost | undefined;
private readonly subagentScheduler: PiSubagentScheduler | undefined;
private readonly interactions: PiInteractionStore;
private readonly extensionUi: PiExtensionUiProjector;
private readonly onExtensionUiProjection: ((projection: PiExtensionUiProjection) => void) | undefined;
private readonly acquireBackgroundLease:
| ((lease: { id: string; kind: 'coding-run' }) => () => void)
| undefined;
private readonly runBackgroundLeases = new Map<
string,
{ runId: string; release: () => void }
>();
private readonly states = new Map<string, ConversationReducerState>();
private readonly inputs = new Map<string, PrepareConversationInput>();
private readonly listeners = new Set<(patch: ConversationPatchEnvelope) => void>();
private readonly projectors = new Map<string, PiEventProjector>();
private readonly projectionChains = new Map<string, Promise<void>>();
private readonly hydrationFlights = new Map<
string,
{ generation: number; flight: Promise<void> }
>();
private readonly unsubscribePool: () => void;
constructor(options: PiConversationRuntimeOptions) {
this.pool = options.pool;
this.registry = options.registry;
this.resolveModel = options.resolveModel;
this.resolveImages = options.resolveImages ?? (async (attachments) => {
if (attachments.length > 0) {
throw new CodingRuntimeContractError(
'CODING_RUNTIME_START_FAILED',
'Conversation attachments are not ready',
true,
);
}
return [];
});
this.projectImage = options.projectImage;
this.createRuntimeId = options.createId ?? (() => randomUUID());
this.now = options.now ?? Date.now;
this.providerRefresh = options.providerRefreshCoordinator ?? new PiProviderRefreshCoordinator();
this.isAuthenticationError = options.isAuthenticationError;
this.refreshCredential = options.refreshCredential;
this.extensionHost = options.extensionHost;
this.subagentScheduler = options.subagentScheduler;
if (options.subagentScheduler && !this.extensionHost) {
throw new Error('Subagent scheduler requires the managed extension host');
}
if (options.subagentScheduler) {
this.extensionHost?.configureSubagents({
scheduler: options.subagentScheduler,
trackGenerationResource: (input) => this.pool.trackGenerationResource(input),
});
}
this.interactions = new PiInteractionStore(this.pool, (interaction) => {
this.emit(interaction.conversationId, {
op: 'interaction.remove',
interactionId: interaction.id,
}, interaction.runId);
});
this.extensionUi = new PiExtensionUiProjector({
getDraftRevision: options.getDraftRevision ?? (() => 0),
...(options.knownExtensionWidgetKeys
? { knownWidgetKeys: options.knownExtensionWidgetKeys }
: {}),
});
this.onExtensionUiProjection = options.onExtensionUiProjection;
this.acquireBackgroundLease = options.acquireBackgroundLease;
if (Boolean(this.isAuthenticationError) !== Boolean(this.refreshCredential)) {
throw new Error('Provider authentication detection and refresh must be configured together');
}
this.unsubscribePool = this.pool.subscribe((event) => this.onPoolEvent(event));
}
async prepare(input: PrepareConversationInput): Promise<ConversationRuntimeState> {
const registered = await this.registry.prepare(input);
const canonicalInput: PrepareConversationInput = {
conversationId: registered.conversation.id,
projectId: input.projectId,
agentId: registered.conversation.agentId,
title: registered.conversation.title,
model: {
model: registered.conversation.model ? clone(registered.conversation.model) : null,
modelResolution: registered.conversation.modelResolution,
},
};
const worker = await this.prepareWorker(canonicalInput);
await this.registry.ensureBinding(canonicalInput, async () => clone(worker.session));
this.inputs.set(input.conversationId, clone(canonicalInput));
const isNewState = !this.states.has(input.conversationId);
if (isNewState) {
this.states.set(input.conversationId, createConversationReducerState(
emptySnapshot(canonicalInput, worker),
));
this.resetProjector(input.conversationId);
} else {
const nextWorker = publicWorkerState(worker);
if (!samePublicWorkerState(this.snapshot(input.conversationId).worker, nextWorker)) {
this.emit(input.conversationId, { op: 'worker.state', state: nextWorker });
}
}
if (isNewState) {
try {
await this.requestHydration(input.conversationId, worker, false);
} catch (error) {
this.recordProjectionFailure(input.conversationId, worker.generation, error);
throw error;
}
}
return this.runtimeState(input.conversationId);
}
async getSnapshot(conversationId: string): Promise<ConversationSnapshot> {
await this.waitForProjection(conversationId);
return clone(this.snapshot(conversationId));
}
async prompt(input: PromptConversationInput): Promise<PromptAcceptance> {
await this.waitForProjection(input.conversationId);
this.assertNoUncertainMutation(input.conversationId);
if (input.mode === 'steer') {
const acceptance = await this.steer(input);
return {
accepted: true,
conversationId: input.conversationId,
clientRequestId: input.clientRequestId,
runId: this.snapshot(input.conversationId).run.runId ?? this.id('run'),
mode: 'steer',
queuePosition: acceptance.queuePosition,
};
}
if (input.mode === 'follow-up') {
const acceptance = await this.followUp(input);
return {
accepted: true,
conversationId: input.conversationId,
clientRequestId: input.clientRequestId,
runId: this.snapshot(input.conversationId).run.runId ?? this.id('run'),
mode: 'follow-up',
queuePosition: acceptance.queuePosition,
};
}
this.snapshot(input.conversationId);
const managedRequest = await this.prepareManagedRequest(input.conversationId, input.attachments.length > 0);
const images = await this.resolveImages(input.attachments);
const runId = this.id('run');
this.acquireRunBackgroundLease(input.conversationId, runId);
const messageId = `client:${input.clientRequestId}`;
this.emit(input.conversationId, {
op: 'message.upsert',
node: {
kind: 'message',
id: messageId,
clientRequestId: input.clientRequestId,
role: 'user',
status: 'optimistic',
blocks: [
...(input.text.length > 0
? [{ kind: 'text' as const, id: `${messageId}:content:0`, text: input.text, status: 'complete' as const }]
: []),
...input.attachments.map(({ attachmentId }, index) => ({
kind: 'image' as const,
id: `${messageId}:content:${index + (input.text.length > 0 ? 1 : 0)}`,
attachmentId,
mime: images[index].mimeType,
})),
],
},
});
const command: PiRpcCommand = {
type: 'prompt',
message: input.text,
...(images.length > 0 ? { images } : {}),
};
const generation = this.pool.getState(input.conversationId)?.generation;
if (generation) this.extensionUi.beginRun(input.conversationId, generation, runId);
let ticket;
try {
if (this.extensionHost && generation) {
await this.extensionHost.bindRun(input.conversationId, generation, runId, managedRequest);
}
ticket = this.pool.startTopLevel({
conversationId: input.conversationId,
runId,
command,
});
} catch (error) {
this.extensionUi.endRun(input.conversationId, runId);
if (this.extensionHost && generation) {
await this.extensionHost.clearRun(input.conversationId, generation, runId).catch(() => undefined);
}
this.releaseRunBackgroundLease(input.conversationId, runId);
throw error;
}
this.emit(input.conversationId, {
op: 'run.state',
run: {
status: ticket.queuePosition ? 'queued' : 'running',
runId,
mode: 'prompt',
startedAt: this.now(),
},
}, runId);
const acceptance = this.acceptPrompt(input.conversationId, runId, command, ticket.accepted);
if (ticket.queuePosition) void acceptance.catch(() => undefined);
else await acceptance;
return {
accepted: true,
conversationId: input.conversationId,
clientRequestId: input.clientRequestId,
runId,
mode: 'prompt',
...(ticket.queuePosition ? { queuePosition: ticket.queuePosition } : {}),
};
}
async steer(input: QueueMessageInput): Promise<QueueAcceptance> {
return await this.queue('steer', input);
}
async followUp(input: QueueMessageInput): Promise<QueueAcceptance> {
return await this.queue('follow-up', input);
}
async abort(conversationId: string): Promise<void> {
await this.waitForProjection(conversationId);
const current = this.snapshot(conversationId).run;
if (runIsTerminal(current.status)) return;
const generation = this.snapshot(conversationId).cursor.workerGeneration;
this.emit(conversationId, {
op: 'run.state',
run: { ...current, status: 'aborting' },
}, current.runId);
if (current.runId) await this.interactions.cancelRun(conversationId, current.runId, true);
try {
await this.pool.request(conversationId, { type: 'abort' });
const generation = this.pool.getState(conversationId)?.generation;
if (current.runId && generation) {
await this.extensionHost?.clearRun(conversationId, generation, current.runId);
this.extensionUi.endRun(conversationId, current.runId);
}
} catch (error) {
await this.waitForProjection(conversationId);
const latest = this.snapshot(conversationId).run;
const worker = this.pool.getState(conversationId);
if (latest.runId === current.runId
&& latest.status === 'aborting'
&& (!worker || worker.state === 'crashed' || worker.generation !== generation)) {
await this.failRun(conversationId, current.runId!, error, generation);
} else if (latest.runId === current.runId && latest.status === 'aborting') {
this.emit(conversationId, { op: 'run.state', run: current }, current.runId);
}
throw error;
}
}
async validateModel(model: ProductModelRef): Promise<ProductModelRef> {
const selection = await this.resolveModel(model);
if (selection.managedCapability) validateManagedReasoningChoice(model.reasoningChoice ?? { mode: 'default' }, selection.managedCapability);
return {
accountId: selection.accountId,
modelId: selection.modelId,
thinkingLevel: model.thinkingLevel,
...(selection.managedCapability ? { reasoningChoice: model.reasoningChoice ?? { mode: 'default' as const } } : {}),
};
}
async setModel(input: SetConversationModelInput): Promise<ConversationModelState> {
await this.waitForProjection(input.conversationId);
this.assertNoUncertainMutation(input.conversationId);
const snapshot = this.snapshot(input.conversationId);
const thinkingLevel = snapshot.conversation.model.model?.thinkingLevel ?? 'off';
const selection = await this.resolveModel({
accountId: input.accountId,
modelId: input.modelId,
thinkingLevel,
});
const model: ConversationModelState = {
model: {
accountId: selection.accountId,
modelId: selection.modelId,
thinkingLevel,
...(selection.managedCapability ? { reasoningChoice: { mode: 'default' as const } } : {}),
},
modelResolution: 'resolved',
};
if (snapshot.conversation.model.model?.accountId !== selection.accountId
|| this.pool.getState(input.conversationId)?.state === 'crashed') {
const persisted = await this.registry.setModel(input.conversationId, model);
this.replaceModel(input.conversationId, persisted);
try {
const worker = await this.pool.reconfigureConversationModel(input.conversationId, persisted);
if (worker) await this.requestHydration(input.conversationId, worker, true);
} catch (error) {
const worker = this.pool.getState(input.conversationId);
if (worker) {
if (this.snapshot(input.conversationId).cursor.workerGeneration !== worker.generation) {
this.replaceWorkerGeneration(input.conversationId, worker, true);
}
this.recordProjectionFailure(input.conversationId, worker.generation, error);
}
throw error;
}
return clone(this.snapshot(input.conversationId).conversation.model);
}
await this.pool.request(input.conversationId, {
type: 'set_model',
provider: selection.runtimeProviderId,
modelId: selection.modelId,
});
const capabilities = await this.pool.request<{ levels?: unknown }>(
input.conversationId,
{ type: 'get_available_thinking_levels' },
{ retry: 'read-only-once' },
);
const state = await this.pool.request(
input.conversationId,
{ type: 'get_state' },
{ retry: 'read-only-once' },
);
return await this.persistEffectiveThinking(
input.conversationId,
model,
state.data,
capabilities.data,
true,
true,
);
}
async setThinking(input: SetThinkingLevelInput): Promise<ConversationModelState> {
await this.waitForProjection(input.conversationId);
this.assertNoUncertainMutation(input.conversationId);
const current = this.snapshot(input.conversationId).conversation.model;
if (!current.model) {
throw new CodingRuntimeContractError(
'CODING_MIGRATION_MODEL_REQUIRED',
'Conversation model must be selected first',
true,
);
}
const selection = current.model.accountId === 'niancode-user-models' ? await this.resolveModel(current.model) : undefined;
if (selection?.managedCapability) {
const choice = input.reasoningChoice ?? (input.thinkingLevel === 'off'
? { mode: 'disabled' as const } : { mode: 'enabled' as const, effort: input.thinkingLevel });
validateManagedReasoningChoice(choice, selection.managedCapability);
await this.pool.request(input.conversationId, { type: 'set_thinking_level', level: managedPiThinkingLevel(choice) });
return this.persistEffectiveThinking(input.conversationId, {
model: { ...current.model, thinkingLevel: 'off', reasoningChoice: choice }, modelResolution: 'resolved',
}, {}, {}, true, true);
}
const capabilities = await this.pool.request<{ levels?: unknown }>(
input.conversationId,
{ type: 'get_available_thinking_levels' },
{ retry: 'read-only-once' },
);
const available = availableThinkingLevels(capabilities.data);
const withCapabilities: ConversationModelState = {
...clone(current),
...(available.length > 0 ? { availableThinkingLevels: available } : {}),
};
this.replaceModel(input.conversationId, withCapabilities);
if (!available.includes(input.thinkingLevel)) {
throw new CodingRuntimeContractError(
'CODING_MODEL_UNAVAILABLE',
'The selected thinking level is not supported by this model',
true,
);
}
await this.pool.request(input.conversationId, {
type: 'set_thinking_level',
level: input.thinkingLevel,
});
const state = await this.pool.request(
input.conversationId,
{ type: 'get_state' },
{ retry: 'read-only-once' },
);
const persisted = await this.persistEffectiveThinking(input.conversationId, {
model: { ...current.model, thinkingLevel: input.thinkingLevel },
modelResolution: 'resolved',
}, state.data, capabilities.data, true, true);
if (persisted.model?.thinkingLevel !== input.thinkingLevel) {
throw new CodingRuntimeContractError(
'CODING_MODEL_UNAVAILABLE',
'Pi did not accept the selected thinking level',
true,
);
}
return clone(persisted);
}
async compact(conversationId: string): Promise<void> {
await this.waitForProjection(conversationId);
this.assertNoUncertainMutation(conversationId);
const managedRequest = await this.prepareManagedRequest(conversationId, false);
const runId = this.id('run');
this.acquireRunBackgroundLease(conversationId, runId);
const generation = this.pool.getState(conversationId)?.generation;
if (generation) this.extensionUi.beginRun(conversationId, generation, runId);
let ticket;
try {
if (this.extensionHost && generation) {
await this.extensionHost.bindRun(conversationId, generation, runId, managedRequest);
}
ticket = this.pool.startTopLevel({
conversationId,
runId,
command: { type: 'compact' },
});
} catch (error) {
this.extensionUi.endRun(conversationId, runId);
if (this.extensionHost && generation) {
await this.extensionHost.clearRun(conversationId, generation, runId).catch(() => undefined);
}
this.releaseRunBackgroundLease(conversationId, runId);
throw error;
}
this.emit(conversationId, {
op: 'run.state',
run: { status: 'compacting', runId, startedAt: this.now() },
}, runId);
try {
await ticket.accepted;
} catch (error) {
if (error instanceof PiProcessError && error.code === 'PI_RPC_TIMEOUT') {
const failure = requestUncertainError();
this.markRunUncertain(conversationId, runId, failure.publicError);
throw failure;
}
await this.failRun(conversationId, runId, error);
throw error;
}
}
async fork(input: ForkConversationInput): Promise<ForkResult> {
await this.waitForProjection(input.sourceConversationId);
this.assertNoUncertainMutation(input.sourceConversationId);
this.snapshot(input.sourceConversationId);
const registered = await this.registry.prepare(input.conversation);
const canonicalInput: PrepareConversationInput = {
conversationId: registered.conversation.id,
projectId: input.conversation.projectId,
agentId: registered.conversation.agentId,
title: registered.conversation.title,
model: {
model: registered.conversation.model ? clone(registered.conversation.model) : null,
modelResolution: registered.conversation.modelResolution,
},
};
const worker = await this.pool.fork(
input.sourceConversationId,
canonicalInput,
input.sourceEntryId,
);
await this.registry.ensureBinding(canonicalInput, async () => clone(worker.session));
this.inputs.set(canonicalInput.conversationId, clone(canonicalInput));
const snapshot = emptySnapshot(canonicalInput, worker);
this.states.set(canonicalInput.conversationId, createConversationReducerState(snapshot));
this.resetProjector(canonicalInput.conversationId);
try {
await this.requestHydration(canonicalInput.conversationId, worker, false);
} catch (error) {
this.recordProjectionFailure(canonicalInput.conversationId, worker.generation, error);
throw error;
}
return {
conversationId: canonicalInput.conversationId,
snapshot: clone(this.snapshot(canonicalInput.conversationId)),
};
}
async recover(conversationId: string): Promise<ConversationRuntimeState> {
await this.waitForProjection(conversationId);
const before = this.snapshot(conversationId);
const generation = this.pool.getState(conversationId)?.generation;
try {
if (before.run.runId && generation) {
await this.interactions.cancelRun(conversationId, before.run.runId, true);
this.extensionUi.endRun(conversationId, before.run.runId);
await this.extensionHost?.clearRun(conversationId, generation, before.run.runId);
}
const state = await this.pool.recover(conversationId);
await this.requestHydration(conversationId, state, false);
this.settleRecoveredRun(conversationId);
return this.runtimeState(conversationId);
} finally {
if (before.run.runId) {
this.releaseRunBackgroundLease(conversationId, before.run.runId);
}
}
}
async dispose(
conversationId: string,
reason: CodingRuntimeDisposeReason,
): Promise<void> {
await this.waitForProjection(conversationId);
const before = this.states.get(conversationId)?.snapshot;
if (reason === 'background_sleep'
&& (this.runBackgroundLeases.has(conversationId)
|| Boolean(before && !runIsTerminal(before.run.status)))) return;
const runId = before?.run.runId;
if (runId && before && !runIsTerminal(before.run.status)) {
await this.enqueueProjection(conversationId, async () => {
const current = this.states.get(conversationId)?.snapshot;
if (!current || current.run.runId !== runId || runIsTerminal(current.run.status)) return;
await this.failRun(
conversationId,
runId,
new CodingRuntimeContractError(
'CODING_RUNTIME_START_FAILED',
'本地 Agent 已中断,原请求未自动重发。',
true,
),
current.cursor.workerGeneration,
false,
);
});
}
const state = this.pool.getState(conversationId);
let stopped = reason !== 'background_sleep';
try {
if (state) await this.interactions.cancelGeneration(conversationId, state.generation);
if (reason === 'background_sleep' && this.hasConversationActiveWork(conversationId)) return;
const didStop = await this.pool.dispose(
conversationId,
reason,
reason === 'background_sleep'
? () => !this.hasConversationActiveWork(conversationId)
: () => true,
);
stopped = reason === 'background_sleep' ? didStop : true;
} finally {
if (stopped) this.releaseConversationBackgroundLease(conversationId);
}
if (reason === 'background_sleep' && !stopped) return;
this.registry.forget(conversationId);
this.inputs.delete(conversationId);
this.states.delete(conversationId);
this.projectors.delete(conversationId);
this.projectionChains.delete(conversationId);
this.hydrationFlights.delete(conversationId);
}
async listCommands(conversationId: string): Promise<CodingRuntimeCommand[]> {
const state = this.pool.getState(conversationId);
if (!state || state.state === 'spawning' || state.state === 'crashed') return [];
const response = await this.pool.request<unknown>(
conversationId,
{ type: 'get_commands' },
{ retry: 'read-only-once' },
);
const record = response.data && typeof response.data === 'object' && !Array.isArray(response.data)
? response.data as Record<string, unknown>
: null;
const candidates = Array.isArray(response.data)
? response.data
: Array.isArray(record?.commands)
? record.commands
: [];
return candidates.flatMap((candidate) => {
if (!candidate || typeof candidate !== 'object' || Array.isArray(candidate)) return [];
const command = candidate as Record<string, unknown>;
const name = typeof command.name === 'string' ? command.name.trim() : '';
if (!name) return [];
return [{
name,
...(typeof command.description === 'string'
? { description: command.description }
: {}),
}];
});
}
async listInteractions(conversationId?: string): Promise<ConversationInteraction[]> {
return this.interactions.list(conversationId);
}
getDiagnostics(): CodingRuntimeDiagnostics {
return this.pool.getDiagnostics();
}
async injectWorkerFailureForProof(
conversationId: string,
failure: PiWorkerProofFailure,
): Promise<{ generation: number }> {
return await this.pool.injectFailureForProof(conversationId, failure);
}
delayNextTopLevelConfirmationForProof(
conversationId: string,
commandType: 'prompt' | 'compact',
delayMs: number,
): void {
this.pool.delayNextTopLevelConfirmationForProof(conversationId, commandType, delayMs);
}
getResilienceProofDiagnostics(): {
pool: ReturnType<PiWorkerPool['getResilienceProofDiagnostics']>;
subagents: ReturnType<PiSubagentScheduler['getDiagnostics']> | null;
extension: ReturnType<PiManagedExtensionHost['getDiagnostics']> | null;
backgroundLeases: { active: number };
} {
return {
pool: this.pool.getResilienceProofDiagnostics(),
subagents: this.subagentScheduler?.getDiagnostics() ?? null,
extension: this.extensionHost?.getDiagnostics() ?? null,
backgroundLeases: { active: this.runBackgroundLeases.size },
};
}
hasActiveWork(): boolean {
if (this.runBackgroundLeases.size > 0) return true;
return [...this.states.values()].some(({ snapshot }) => !runIsTerminal(snapshot.run.status));
}
private hasConversationActiveWork(conversationId: string): boolean {
if (this.runBackgroundLeases.has(conversationId)) return true;
const snapshot = this.states.get(conversationId)?.snapshot;
return Boolean(snapshot && !runIsTerminal(snapshot.run.status));
}
private assertNoUncertainMutation(conversationId: string): void {
const run = this.states.get(conversationId)?.snapshot?.run;
if (!runIsTerminal(run?.status ?? 'idle')
&& run?.error?.code === 'CODING_REQUEST_UNCERTAIN') {
throw requestUncertainError();
}
}
markProviderStale(): void {
this.pool.markProviderStale();
}
markResourcesStale(): void {
this.pool.markResourcesStale();
}
async refreshResources(): Promise<void> {
await this.pool.refreshResources();
}
subscribe(listener: (patch: ConversationPatchEnvelope) => void): () => void {
this.listeners.add(listener);
return () => this.listeners.delete(listener);
}
async respondInteraction(
conversationId: string,
response: ConversationInteractionResponse,
): Promise<void> {
await this.interactions.respond(conversationId, response);
}
async shutdown(): Promise<void> {
this.unsubscribePool();
try {
const results = await Promise.allSettled([
this.pool.shutdown(),
this.extensionHost?.close(),
]);
const failure = results.find((result): result is PromiseRejectedResult => (
result.status === 'rejected'
));
if (failure) throw failure.reason;
} finally {
for (const conversationId of [...this.runBackgroundLeases.keys()]) {
this.releaseConversationBackgroundLease(conversationId);
}
this.projectors.clear();
this.projectionChains.clear();
this.hydrationFlights.clear();
}
}
private async queue(
mode: 'steer' | 'follow-up',
input: QueueMessageInput,
): Promise<QueueAcceptance> {
await this.waitForProjection(input.conversationId);
this.assertNoUncertainMutation(input.conversationId);
const snapshot = this.snapshot(input.conversationId);
if (snapshot.run.runId) {
this.acquireRunBackgroundLease(input.conversationId, snapshot.run.runId);
}
await this.prepareManagedRequest(input.conversationId, input.attachments.length > 0);
const images = await this.resolveImages(input.attachments);
const queuePosition = snapshot.queue.items.length + 1;
const queueId = this.id('queue');
this.emit(input.conversationId, {
op: 'queue.replace',
queue: {
items: [
...snapshot.queue.items,
{
id: queueId,
clientRequestId: input.clientRequestId,
mode,
text: input.text,
attachmentIds: input.attachments.map(({ attachmentId }) => attachmentId),
},
],
},
}, snapshot.run.runId);
try {
await this.pool.request(input.conversationId, {
type: mode === 'steer' ? 'steer' : 'follow_up',
message: input.text,
...(images.length > 0 ? { images } : {}),
});
} catch (error) {
const current = this.snapshot(input.conversationId);
this.emit(input.conversationId, {
op: 'queue.replace',
queue: { items: current.queue.items.filter(({ id }) => id !== queueId) },
}, current.run.runId);
throw error;
}
return {
accepted: true,
conversationId: input.conversationId,
clientRequestId: input.clientRequestId,
mode,
queuePosition,
};
}
private async prepareWorker(input: PrepareConversationInput): Promise<PiWorkerPoolState> {
if (!this.isAuthenticationError || !this.refreshCredential || !input.model.model) {
return await this.pool.prepare(input);
}
try {
return await this.providerRefresh.withSingleAuthRecovery({
accountId: input.model.model.accountId,
operation: async () => await this.pool.prepare(input),
isAuthenticationError: this.isAuthenticationError,
refreshCredential: async () => await this.refreshCredential!(input.model.model!.accountId),
reopenWorker: async () => {
if (!this.pool.getState(input.conversationId)) return;
const recovered = await this.pool.recover(input.conversationId);
if (this.states.has(input.conversationId)) {
await this.requestHydration(input.conversationId, recovered, false);
}
},
});
} catch (error) {
if (this.isAuthenticationError(error)) {
throw new CodingRuntimeContractError(
'CODING_PROVIDER_AUTH_REQUIRED',
'Provider authentication failed after one recovery attempt',
true,
);
}
throw error;
}
}
private async acceptPrompt(
conversationId: string,
runId: string,
command: PiRpcCommand,
firstAcceptance: Promise<PiRpcResponse>,
): Promise<void> {
const accountId = this.snapshot(conversationId).conversation.model.model?.accountId;
try {
if (accountId && this.isAuthenticationError && this.refreshCredential) {
let acceptance = firstAcceptance;
await this.providerRefresh.withSingleAuthRecovery({
accountId,
operation: async (attempt) => {
if (attempt === 1) {
acceptance = this.pool.startTopLevel({ conversationId, runId, command }).accepted;
}
return await acceptance;
},
isAuthenticationError: this.isAuthenticationError,
refreshCredential: async () => await this.refreshCredential!(accountId),
reopenWorker: async () => {
const recovered = await this.pool.recover(conversationId);
await this.requestHydration(conversationId, recovered, true);
},
});
} else {
await firstAcceptance;
}
const current = this.states.get(conversationId)?.snapshot?.run;
if (current?.runId === runId && current.status === 'queued') {
this.emit(conversationId, {
op: 'run.state',
run: { ...current, status: 'running' },
}, runId);
}
} catch (error) {
const failure = this.isAuthenticationError?.(error)
? new CodingRuntimeContractError(
'CODING_PROVIDER_AUTH_REQUIRED',
'Provider authentication failed after one recovery attempt',
true,
)
: error;
if (failure instanceof PiProcessError && failure.code === 'PI_RPC_TIMEOUT') {
const uncertain = requestUncertainError();
this.markRunUncertain(conversationId, runId, uncertain.publicError);
throw uncertain;
}
this.pool.failTopLevel(
conversationId,
runId,
failure instanceof Error ? failure : new Error('Prompt acceptance failed'),
);
await this.failRun(conversationId, runId, failure);
throw failure;
}
}
private markRunUncertain(
conversationId: string,
runId: string,
error: CodingRuntimePublicError,
): void {
const current = this.states.get(conversationId)?.snapshot?.run;
if (!current || current.runId !== runId || runIsTerminal(current.status)) return;
this.emit(conversationId, {
op: 'run.state',
run: { ...current, error: clone(error) },
}, runId);
}
private snapshot(conversationId: string): ConversationSnapshot {
const snapshot = this.states.get(conversationId)?.snapshot;
if (!snapshot) {
throw new CodingRuntimeContractError(
'CODING_CONVERSATION_NOT_FOUND',
'Conversation is not prepared',
true,
);
}
return snapshot;
}
private runtimeState(conversationId: string): ConversationRuntimeState {
const snapshot = this.snapshot(conversationId);
return {
conversationId,
status: snapshot.worker.status,
workerGeneration: snapshot.worker.generation,
...(snapshot.worker.error ? { error: clone(snapshot.worker.error) } : {}),
};
}
private emit(conversationId: string, patch: ConversationPatch, runId?: string): void {
const state = this.states.get(conversationId);
const snapshot = state?.snapshot;
if (!state || !snapshot) return;
const generation = this.pool.getState(conversationId)?.generation
?? snapshot.cursor.workerGeneration;
const envelope: ConversationPatchEnvelope = {
conversationId,
workerGeneration: generation,
...(runId ? { runId } : {}),
seq: snapshot.cursor.seq + 1,
at: this.now(),
patch: clone(patch),
};
const next = reduceConversationPatch(state, envelope);
if (next.invalidation) {
throw new CodingRuntimeContractError(
'CODING_RUNTIME_PROTOCOL_ERROR',
next.invalidation.reason,
true,
);
}
this.states.set(conversationId, next);
for (const listener of this.listeners) listener(clone(envelope));
}
private onPoolEvent(event: PiWorkerPoolEvent): void {
if (event.type === 'top-level.confirmed') {
void this.enqueueProjection(event.conversationId, async () => {
const snapshot = this.states.get(event.conversationId)?.snapshot;
const current = snapshot?.run;
if (!snapshot
|| snapshot.cursor.workerGeneration !== event.generation
|| current?.runId !== event.runId
|| runIsTerminal(current.status)) return;
const { error: _error, ...confirmed } = current;
this.emit(event.conversationId, {
op: 'run.state',
run: {
...confirmed,
status: confirmed.status === 'queued' ? 'running' : confirmed.status,
},
}, event.runId);
}).catch((error) => {
this.recordProjectionFailure(event.conversationId, event.generation, error);
});
return;
}
if (event.type === 'top-level.settled') {
void this.enqueueProjection(event.conversationId, async () => {
const snapshot = this.states.get(event.conversationId)?.snapshot;
const current = snapshot?.run;
if (!snapshot
|| snapshot.cursor.workerGeneration !== event.generation
|| current?.runId !== event.runId) return;
if (!runIsTerminal(current.status)) {
this.emit(event.conversationId, {
op: 'run.state',
run: {
status: 'idle',
runId: event.runId,
...(current.mode ? { mode: current.mode } : {}),
...(current.startedAt !== undefined ? { startedAt: current.startedAt } : {}),
settledAt: this.now(),
terminalReason: current.status === 'aborting' ? 'aborted' : 'completed',
},
}, event.runId);
}
this.extensionUi.endRun(event.conversationId, event.runId);
try {
await Promise.allSettled([
this.interactions.cancelRun(event.conversationId, event.runId, true),
this.extensionHost?.clearRun(
event.conversationId,
event.generation,
event.runId,
),
]);
} finally {
this.releaseRunBackgroundLease(event.conversationId, event.runId);
}
if (event.source === 'state_probe') {
const worker = this.pool.getState(event.conversationId);
if (worker?.generation === event.generation && worker.state !== 'crashed') {
await this.hydrateGenerationNow(event.conversationId, worker, false);
}
}
}).catch((error) => {
this.recordProjectionFailure(event.conversationId, event.generation, error);
});
return;
}
if (event.type === 'top-level.failed') {
void this.enqueueProjection(event.conversationId, async () => {
await this.failRun(
event.conversationId,
event.runId,
event.error,
event.generation,
);
}).catch((error) => {
this.recordProjectionFailure(event.conversationId, event.generation, error);
});
return;
}
if (event.type === 'worker.replaced') {
if (!this.states.has(event.conversationId)) return;
this.replaceWorkerGeneration(event.conversationId, event.state, true);
this.resetProjector(event.conversationId);
const runId = this.states.get(event.conversationId)?.snapshot.run.runId;
const activeRun = this.pool.getActiveRun(event.conversationId);
const continuesActiveRun = Boolean(runId && activeRun?.runId === runId);
if (runId && continuesActiveRun) {
this.extensionUi.replaceGeneration(event.conversationId, event.generation, runId);
}
if (this.extensionHost && runId && continuesActiveRun) {
void this.extensionHost.bindRun(event.conversationId, event.generation, runId).catch((error) => {
this.recordProjectionFailure(event.conversationId, event.generation, error);
});
}
void this.requestHydration(event.conversationId, event.state, true).catch((error) => {
this.recordProjectionFailure(event.conversationId, event.generation, error);
});
return;
}
if (event.type === 'worker.crashed') {
void this.enqueueProjection(event.conversationId, async () => {
const snapshot = this.states.get(event.conversationId)?.snapshot;
if (!snapshot || snapshot.cursor.workerGeneration !== event.generation) return;
await this.interactions.cancelGeneration(event.conversationId, event.generation);
const current = this.snapshot(event.conversationId).run;
if (current.runId && !runIsTerminal(current.status)) {
await this.failRun(
event.conversationId,
current.runId,
event.error,
event.generation,
);
}
const state = this.pool.getState(event.conversationId);
if (state?.generation === event.generation) {
this.emit(event.conversationId, {
op: 'worker.state',
state: publicWorkerState(state),
});
}
}).catch((error) => {
this.recordProjectionFailure(event.conversationId, event.generation, error);
});
return;
}
void this.enqueueProjection(event.conversationId, async () => {
const snapshot = this.states.get(event.conversationId)?.snapshot;
if (!snapshot || snapshot.cursor.workerGeneration !== event.generation) return;
if (event.event.type === 'extension_ui_request' && snapshot.run.runId) {
const interaction = this.interactions.open(
event.conversationId,
event.generation,
snapshot.run.runId,
event.event,
);
if (interaction) {
this.emit(event.conversationId, { op: 'interaction.upsert', interaction }, interaction.runId);
return;
}
const projection = this.extensionUi.project(
event.conversationId,
event.generation,
snapshot.run.runId,
event.event,
);
if (projection) this.onExtensionUiProjection?.(projection);
return;
}
const projector = this.projector(event.conversationId);
const patches = await projector.project(snapshot, event.event);
for (const patch of patches) {
this.emit(event.conversationId, patch, this.snapshot(event.conversationId).run.runId);
}
if (event.event.type === 'agent_end' || event.event.type === 'agent_settled') {
const worker = this.pool.getState(event.conversationId);
if (worker?.generation === event.generation && worker.state !== 'crashed') {
await this.hydrateGenerationNow(
event.conversationId,
worker,
event.event.type === 'agent_end',
);
}
}
if (event.event.type === 'agent_settled' && snapshot.run.runId) {
try {
await Promise.allSettled([
this.interactions.cancelRun(event.conversationId, snapshot.run.runId, true),
this.extensionHost?.clearRun(
event.conversationId,
event.generation,
snapshot.run.runId,
),
]);
this.extensionUi.endRun(event.conversationId, snapshot.run.runId);
} finally {
this.releaseRunBackgroundLease(event.conversationId, snapshot.run.runId);
}
}
}).catch((error) => {
this.recordProjectionFailure(event.conversationId, event.generation, error);
});
}
private async failRun(
conversationId: string,
runId: string,
error: unknown,
generation?: number,
releaseBackgroundLease = true,
): Promise<void> {
const current = this.states.get(conversationId)?.snapshot?.run;
const snapshotGeneration = this.states.get(conversationId)?.snapshot?.cursor.workerGeneration;
if (current?.runId !== runId
|| (generation !== undefined && snapshotGeneration !== generation)) return;
if (runIsTerminal(current.status)) {
if (releaseBackgroundLease) this.releaseRunBackgroundLease(conversationId, runId);
return;
}
const publicError = runtimeFailure(error);
this.emit(conversationId, {
op: 'run.state',
run: {
status: 'error',
runId,
...(current.mode ? { mode: current.mode } : {}),
...(current.startedAt !== undefined ? { startedAt: current.startedAt } : {}),
settledAt: this.now(),
terminalReason: 'failed',
error: publicError,
},
}, runId);
logger.warn('[PiWorkerLifecycle]', {
event: 'run.failed',
conversationId,
runId,
generation: generation ?? snapshotGeneration,
code: publicError.code,
});
this.extensionUi.endRun(conversationId, runId);
const runGeneration = generation ?? this.pool.getState(conversationId)?.generation;
try {
await Promise.allSettled([
this.interactions.cancelRun(conversationId, runId, true),
runGeneration
? this.extensionHost?.clearRun(conversationId, runGeneration, runId)
: undefined,
]);
} finally {
if (releaseBackgroundLease) this.releaseRunBackgroundLease(conversationId, runId);
}
}
private requestHydration(
conversationId: string,
workerState: PiWorkerPoolState,
preserveRun: boolean,
): Promise<void> {
const existing = this.hydrationFlights.get(conversationId);
if (existing?.generation === workerState.generation) return existing.flight;
const hydration = this.enqueueProjection(conversationId, async () => {
await this.hydrateGenerationNow(conversationId, workerState, preserveRun);
});
const flight = hydration.finally(() => {
if (this.hydrationFlights.get(conversationId)?.flight === flight) {
this.hydrationFlights.delete(conversationId);
}
});
this.hydrationFlights.set(conversationId, { generation: workerState.generation, flight });
return flight;
}
private async hydrateGenerationNow(
conversationId: string,
workerState: PiWorkerPoolState,
preserveRun: boolean,
): Promise<void> {
const currentWorker = this.pool.getState(conversationId);
if (!currentWorker || currentWorker.generation !== workerState.generation) return;
if (this.snapshot(conversationId).cursor.workerGeneration !== workerState.generation) {
this.replaceWorkerGeneration(conversationId, workerState, preserveRun);
this.resetProjector(conversationId);
}
const [stateResponse, entriesResponse, statsResponse, capabilitiesResponse] = await Promise.all([
this.pool.request(conversationId, { type: 'get_state' }, { retry: 'read-only-once' }),
this.pool.request(conversationId, { type: 'get_entries' }, { retry: 'read-only-once' }),
this.pool.request(conversationId, { type: 'get_session_stats' }, { retry: 'read-only-once' }),
this.pool.request(conversationId, { type: 'get_available_thinking_levels' }, {
retry: 'read-only-once',
}),
]);
const before = this.snapshot(conversationId);
let projected = await projectPiSessionSnapshot({
snapshot: before,
workerGeneration: workerState.generation,
state: stateResponse.data,
entries: entriesResponse.data,
stats: statsResponse.data,
...(this.projectImage ? { projectImage: this.projectImage } : {}),
});
if (preserveRun) {
projected = {
...projected,
run: clone(before.run),
queue: clone(before.queue),
};
}
projected = {
...projected,
conversation: {
...projected.conversation,
model: await this.persistEffectiveThinking(
conversationId,
projected.conversation.model,
stateResponse.data,
capabilitiesResponse.data,
false,
),
},
};
this.states.set(conversationId, createConversationReducerState(projected));
}
private async prepareManagedRequest(conversationId: string, hasImages: boolean) {
const model = this.snapshot(conversationId).conversation.model.model;
if (!model || model.accountId !== 'niancode-user-models') return undefined;
const selection = await this.resolveModel(model);
if (!selection.managedCapability) return undefined;
if (hasImages && !selection.input.includes('image')) {
throw new CodingRuntimeContractError('CODING_MODEL_UNAVAILABLE', '该模型尚未确认支持图片输入,请选择支持图片的模型', true);
}
return buildManagedModelRequest(model.modelId, model.reasoningChoice ?? { mode: 'default' }, selection.managedCapability);
}
private async persistEffectiveThinking(
conversationId: string,
requested: ConversationModelState,
stateValue: unknown,
capabilitiesValue: unknown,
replaceSnapshot = true,
forcePersist = false,
): Promise<ConversationModelState> {
if (!requested.model) return clone(requested);
const selection = requested.model.accountId === 'niancode-user-models' ? await this.resolveModel(requested.model) : undefined;
if (selection?.managedCapability) {
const durable: ConversationModelState = {
model: { ...requested.model, thinkingLevel: 'off',
reasoningChoice: requested.model.reasoningChoice ?? { mode: 'default' } },
modelResolution: 'resolved',
};
const persisted = await this.registry.setModel(conversationId, durable);
this.pool.updateConversationModel(conversationId, persisted);
const result = { ...persisted, managedCapability: selection.managedCapability };
if (replaceSnapshot) this.replaceModel(conversationId, result);
return clone(result);
}
const effective = effectiveThinkingLevel(stateValue) ?? requested.model.thinkingLevel;
const available = availableThinkingLevels(capabilitiesValue);
if (!available.includes(effective)) available.push(effective);
const durable: ConversationModelState = {
model: { ...requested.model, thinkingLevel: effective },
modelResolution: 'resolved',
};
const persisted = !forcePersist && requested.model.thinkingLevel === effective
? durable
: await this.registry.setModel(conversationId, durable);
this.pool.updateConversationModel(conversationId, persisted);
const result: ConversationModelState = {
...persisted,
...(available.length > 0 ? { availableThinkingLevels: available } : {}),
};
if (replaceSnapshot) this.replaceModel(conversationId, result);
else {
const input = this.inputs.get(conversationId);
if (input) input.model = clone(persisted);
}
return clone(result);
}
private enqueueProjection(conversationId: string, action: () => Promise<void>): Promise<void> {
const previous = this.projectionChains.get(conversationId) ?? Promise.resolve();
const flight = previous.then(action);
const tail = flight.catch(() => undefined);
this.projectionChains.set(conversationId, tail);
void tail.finally(() => {
if (this.projectionChains.get(conversationId) === tail) {
this.projectionChains.delete(conversationId);
}
});
return flight;
}
private async waitForProjection(conversationId: string): Promise<void> {
await this.projectionChains.get(conversationId);
}
private projector(conversationId: string): PiEventProjector {
const existing = this.projectors.get(conversationId);
if (existing) return existing;
return this.resetProjector(conversationId);
}
private resetProjector(conversationId: string): PiEventProjector {
const projector = new PiEventProjector({
createId: randomUUID,
now: this.now,
...(this.projectImage ? { projectImage: this.projectImage } : {}),
});
this.projectors.set(conversationId, projector);
return projector;
}
private recordProjectionFailure(
conversationId: string,
generation: number,
error: unknown,
): void {
const state = this.states.get(conversationId);
const snapshot = state?.snapshot;
if (!state || !snapshot || snapshot.cursor.workerGeneration !== generation) return;
const publicError = runtimeFailure(error);
this.states.set(conversationId, createConversationReducerState({
...clone(snapshot),
worker: { status: 'error', generation, error: publicError },
run: {
...clone(snapshot.run),
status: 'error',
error: publicError,
},
}));
}
private settleRecoveredRun(conversationId: string): void {
const current = this.snapshot(conversationId);
this.states.set(conversationId, createConversationReducerState({
...clone(current),
run: { status: 'idle' },
queue: { items: [] },
pendingInteractions: [],
}));
}
private replaceWorkerGeneration(
conversationId: string,
workerState: PiWorkerPoolState,
preserveRun: boolean,
): void {
const current = this.snapshot(conversationId);
const next: ConversationSnapshot = {
...clone(current),
run: preserveRun ? clone(current.run) : { status: 'idle' },
queue: preserveRun ? clone(current.queue) : { items: [] },
pendingInteractions: [],
worker: publicWorkerState(workerState),
cursor: {
workerGeneration: workerState.generation,
seq: 0,
...(current.cursor.leafEntryId ? { leafEntryId: current.cursor.leafEntryId } : {}),
},
};
this.states.set(conversationId, createConversationReducerState(next));
}
private replaceModel(conversationId: string, model: ConversationModelState): void {
const state = this.states.get(conversationId);
const snapshot = state?.snapshot;
if (!state || !snapshot) return;
this.states.set(conversationId, createConversationReducerState({
...snapshot,
conversation: { ...snapshot.conversation, model: clone(model) },
}));
const input = this.inputs.get(conversationId);
if (input) input.model = clone(model);
}
private acquireRunBackgroundLease(conversationId: string, runId: string): void {
const current = this.runBackgroundLeases.get(conversationId);
if (current?.runId === runId) return;
if (current) {
throw new CodingRuntimeContractError(
'CODING_RUNTIME_START_FAILED',
'Conversation already has an active run',
true,
);
}
const release = this.acquireBackgroundLease?.({
id: `coding-run:${conversationId}:${runId}`,
kind: 'coding-run',
}) ?? (() => undefined);
this.runBackgroundLeases.set(conversationId, { runId, release });
}
private releaseRunBackgroundLease(conversationId: string, runId: string): void {
const current = this.runBackgroundLeases.get(conversationId);
if (!current || current.runId !== runId) return;
this.runBackgroundLeases.delete(conversationId);
current.release();
}
private releaseConversationBackgroundLease(conversationId: string): void {
const current = this.runBackgroundLeases.get(conversationId);
if (current) this.releaseRunBackgroundLease(conversationId, current.runId);
}
private id(kind: RuntimeIdKind): string {
return this.createRuntimeId(kind);
}
}