import { accessSync, constants, existsSync } from "node:fs"; import { basename, join } from "node:path"; import { fileURLToPath } from "node:url"; import { randomBytes } from "node:crypto"; import type { Readable, Writable } from "node:stream"; import type { ChildReplyEnvelope } from "./child-reply-envelope.ts"; import type { AgentSnapshot } from "./tree-controller.ts"; import { SUPERVISOR_CHANNEL_LIMITS } from "./supervisor-channel.ts"; import { ACTIVITY_MAX_TEXT_BYTES, isSafeModelCallFailureReason, isSafeToolOrigin, isValidToolExecutionGeneration, } from "./rpc-bridge-event.ts"; import { LengthPrefixedFrameDecoder } from "./length-prefixed-frame-decoder.ts"; import { isManagedProcessTreeAdapter, type ExitObservation, type ManagedProcessTransport, type ProcessLaunchSpec, type ProcessTreeAdapter, type ProcessTreeHandle, type ResourceObservation, } from "./process-tree-capability.ts"; import { ManagedRpcStartupError, parseManagedRpcStartupDiagnostic, type ManagedRpcStartupDiagnostic, } from "./startup-diagnostic.ts"; export { ManagedRpcStartupError, type ManagedRpcStartupDiagnostic, } from "./startup-diagnostic.ts"; export type ManagedRpcReply = ChildReplyEnvelope; export type ManagedRpcTransportFault = "eof" | "protocol_fault" | "process_exit"; export const MANAGED_RPC_BRIDGE_PROTOCOL = "wj-pi-subagents/managed-rpc/10" as const; /** 只用于节点启动事务的一次性本地认证,不进入公开控制面。 */ export const MANAGED_RPC_BRIDGE_CREDENTIAL_ENV = "WJ_PI_SUBAGENTS_MANAGED_RPC_CREDENTIAL" as const; /** 操作者显式指定桥接 JS 运行时;优先级低于装配选项,高于宿主形态解析。 */ export const BRIDGE_RUNTIME_ENV = "WJ_PI_SUBAGENTS_BRIDGE_RUNTIME" as const; /** 外层桥接 JSON 正文的硬边界。 */ export const MANAGED_RPC_BRIDGE_MAX_FRAME_BYTES = 64 * 1024; /** 单个高层命令的重组边界,覆盖长模板配置和监督启动快照。 */ export const MANAGED_RPC_BRIDGE_MAX_COMMAND_PAYLOAD_BYTES = 512 * 1024; /** Base64URL 分片为外层桥接 JSON 保留协议字段空间。 */ export const MANAGED_RPC_BRIDGE_COMMAND_CHUNK_BYTES = 45 * 1024; /** * 监督字节进入外层桥接前的传输分片上限。46 KiB 经 Base64URL 后仍能放入 * 64 KiB 外层帧;它不是完整 SupervisorChannel 帧的上限。 */ export const MANAGED_RPC_SUPERVISOR_MAX_FRAME_BYTES = 46 * 1024; /** SupervisorChannel 的完整 JSON 正文边界。 */ export const MANAGED_RPC_SUPERVISOR_MAX_BODY_BYTES = SUPERVISOR_CHANNEL_LIMITS.maxFrameBytes; /** 完整监督帧包含四字节长度头,桥接会将它拆成多个 tunnel chunk。 */ export const MANAGED_RPC_SUPERVISOR_MAX_ENCODED_FRAME_BYTES = MANAGED_RPC_SUPERVISOR_MAX_BODY_BYTES + 4; /** 节点身份只在预留成功后传给桥接进程,用于建立独立监督通道。 */ export interface ManagedRpcSupervisorInit { readonly root_id: string; readonly local_agent_id: string; readonly peer_agent_id: string; readonly parent_agent_id: string | null; readonly depth: number; readonly credential: string; readonly initial_snapshot: readonly AgentSnapshot[]; readonly initial_subtree_revision: number; } export interface ManagedRpcNodeStartContext { readonly supervisor?: ManagedRpcSupervisorInit; /** 身份预留后由根运行时快照派生,不能覆盖桥接一次性凭据。 */ readonly environment?: Readonly>; } /** 桥接进程只暴露高层命令,不暴露 Pi JSONL 流。 */ export interface ManagedRpcBridge { start(signal?: AbortSignal, context?: ManagedRpcNodeStartContext): Promise; prompt(message: string): Promise; /** 同一持续会话的自适应发送:Pi active 时 steer,idle 时启动 prompt。 */ steer(message: string): Promise; abort(): Promise; getState(): Promise; requestClose(signal: AbortSignal): Promise; onEvent(listener: (event: unknown) => void): () => void; onTransportFault(listener: (fault: ManagedRpcTransportFault) => void): () => void; /** 与任务 RPC 复用同一读取者的父子监督帧转发。 */ sendSupervisorFrame(frame: Uint8Array): Promise; onSupervisorFrame(listener: (frame: Uint8Array) => void): () => void; release?(): Promise; } export interface ManagedRpcBridgeFactoryOptions { /** 仅父端桥接客户端使用;不得写入日志、回复或树快照。 */ readonly credential?: string; /** 仅通过首个有界 start 命令交给 bridge,不进入任何进程命令行。 */ readonly rpcOptions?: Readonly>; } export type ManagedRpcBridgeFactory = ( transport: ManagedProcessTransport, options?: ManagedRpcBridgeFactoryOptions, ) => ManagedRpcBridge; export interface ManagedRpcNodeLaunchOptions { readonly command: string; readonly args?: readonly string[]; readonly cwd?: string; readonly env?: NodeJS.ProcessEnv; } export interface ManagedRpcNodeOptions { readonly processTreeAdapter: ProcessTreeAdapter; readonly launch: ManagedRpcNodeLaunchOptions; /** * 装配期已确认的启动阻塞。start 在触碰进程树之前立即抛出,让监督器把 * 环境问题归类为 spawn_failed,而不是在装配处被吞成 internal_error。 */ readonly startupBlocked?: Error; /** 首个 bridge start 命令携带的 Pi RpcClient 配置,不进入 OS 命令行。 */ readonly rpcOptions?: Readonly>; /** 测试可注入 bridge;生产默认使用有界本地帧桥。 */ readonly bridgeFactory?: ManagedRpcBridgeFactory; } export interface ManagedRpcNodeAssemblyOptions { readonly processTreeAdapter: ProcessTreeAdapter; readonly cwd?: string; readonly env?: NodeJS.ProcessEnv; readonly rpcOptions?: Readonly>; /** 测试/打包时可指定已编译的桥接入口。 */ readonly bridgeScriptPath?: string; readonly bridgeFactory?: ManagedRpcBridgeFactory; /** 装配级运行时覆盖;优先于环境变量与宿主形态解析。 */ readonly bridgeRuntimePath?: string; /** 测试可替换 PATH 查找输入;生产使用进程环境。 */ readonly bridgeRuntimePathEnv?: string; /** 测试/宿主可覆盖宿主可执行文件路径;生产使用 process.execPath。 */ readonly hostExecPath?: string; } /** 编译宿主上找不到可用 JS 运行时;节点启动时立即抛出,由监督器归类为 spawn_failed。 */ export class BridgeRuntimeUnavailableError extends Error { constructor() { super("未找到可用的桥接 JS 运行时(node/bun)"); this.name = "BridgeRuntimeUnavailableError"; } } /** * 解析运行桥接脚本的 JS 运行环境。通常与宿主相同;但单文件编译宿主 * (execPath 是产品二进制而非 node/bun)不能执行脚本,此时回退到 * PATH 上的 node(其次 bun)。找不到时返回 undefined,由装配入口转为 * 启动期阻塞错误,避免用产品二进制重演启动超时。 */ export function resolveBridgeRuntime( execPath: string = process.execPath, pathEnv: string | undefined = process.env.PATH, platform: string = process.platform, ): string | undefined { const base = basename(execPath).toLowerCase().replace(/\.exe$/, ""); if (base === "node" || base === "nodejs" || base === "bun") return execPath; for (const name of ["node", "bun"]) { const found = findExecutableOnPath(name, pathEnv, platform); if (found !== undefined) return found; } return undefined; } function findExecutableOnPath( name: string, pathEnv: string | undefined, platform: string, ): string | undefined { if (pathEnv === undefined || pathEnv === "") return undefined; const separator = platform === "win32" ? ";" : ":"; // Windows Job Object helper 以 CreateProcess 启动,无法执行 .cmd 或无扩展名 // 脚本;只接受 .exe,避免选中不可启动的候选并阻断后续 bun 回退。 const candidates = platform === "win32" ? [`${name}.exe`] : [name]; for (const dir of pathEnv.split(separator)) { if (dir === "") continue; for (const candidate of candidates) { const full = join(dir, candidate); try { if (!existsSync(full)) continue; if (platform !== "win32") accessSync(full, constants.X_OK); } catch { continue; } return full; } } return undefined; } function needsTypeStrippingFlag(runtimePath: string, scriptPath: string): boolean { if (!scriptPath.endsWith(".ts")) return false; const base = basename(runtimePath).toLowerCase().replace(/\.(exe|cmd)$/, ""); return base === "node" || base === "nodejs"; } /** 生成平台适配器在启动前接收的桥接进程说明。 */ export function createManagedRpcNodeLaunchSpec( options: Pick< ManagedRpcNodeAssemblyOptions, "cwd" | "env" | "bridgeScriptPath" | "bridgeRuntimePath" | "bridgeRuntimePathEnv" | "hostExecPath" >, ): ManagedRpcNodeLaunchOptions { const scriptPath = options.bridgeScriptPath ?? defaultBridgeScriptPath(); const runtime = options.bridgeRuntimePath ?? nonEmptyRuntime(process.env[BRIDGE_RUNTIME_ENV]) ?? resolveBridgeRuntime(options.hostExecPath ?? process.execPath, options.bridgeRuntimePathEnv); if (runtime === undefined) throw new BridgeRuntimeUnavailableError(); return Object.freeze({ command: runtime, args: Object.freeze([ ...(needsTypeStrippingFlag(runtime, scriptPath) ? ["--experimental-strip-types"] : []), scriptPath, ]), ...(options.cwd === undefined ? {} : { cwd: options.cwd }), ...(options.env === undefined ? {} : { env: Object.freeze({ ...options.env }) }), }); } function nonEmptyRuntime(value: string | undefined): string | undefined { return value === undefined || value === "" ? undefined : value; } /** * Node 原生 type stripping 不允许执行 node_modules 内的 .ts 文件。发布包提供 * 编译 bridge 时优先运行它;源码目录未构建时仍保留本地开发的 .ts 回退路径。 */ export function resolveManagedRpcBridgeScriptPath( moduleUrl: string = import.meta.url, fileExists: (path: string) => boolean = existsSync, ): string { // 编译后的模块本身位于 dist/src,优先使用同目录 bridge;源码模块则继续查找 dist。 const colocatedCompiled = fileURLToPath(new URL("./rpc-bridge-process.js", moduleUrl)); if (fileExists(colocatedCompiled)) return colocatedCompiled; const compiled = fileURLToPath(new URL("../dist/src/rpc-bridge-process.js", moduleUrl)); if (fileExists(compiled)) return compiled; return fileURLToPath(new URL("./rpc-bridge-process.ts", moduleUrl)); } function defaultBridgeScriptPath(): string { return resolveManagedRpcBridgeScriptPath(); } /** 生产装配便捷入口;返回值仍是单一 `ManagedRpcNode` 深模块。 */ export function createManagedRpcNode(options: ManagedRpcNodeAssemblyOptions): ManagedRpcNode { let launch: ManagedRpcNodeLaunchOptions; let startupBlocked: Error | undefined; try { launch = createManagedRpcNodeLaunchSpec(options); } catch (error) { if (!(error instanceof BridgeRuntimeUnavailableError)) throw error; // 装配错误在控制器边界会被吞成 internal_error;延迟到 start 抛出, // 让监督器把环境缺失归类为 spawn_failed。占位启动说明不会被使用。 startupBlocked = error; launch = Object.freeze({ command: process.execPath, args: Object.freeze([]) }); } return new ManagedRpcNode({ processTreeAdapter: options.processTreeAdapter, launch, ...(startupBlocked === undefined ? {} : { startupBlocked }), ...(options.rpcOptions === undefined ? {} : { rpcOptions: options.rpcOptions }), ...(options.bridgeFactory === undefined ? {} : { bridgeFactory: options.bridgeFactory }), }); } export interface ManagedRpcNodeLike { readonly process_binding: "managed"; start(signal?: AbortSignal, context?: ManagedRpcNodeStartContext): Promise; prompt(message: string): Promise; /** 与 bridge 一致,steer 必须由宿主在 active/idle 间原子裁决。 */ steer(message: string): Promise; abort(): Promise; getState(): Promise; onEvent(listener: (event: unknown) => void): () => void; onTransportFault(listener: (fault: ManagedRpcTransportFault) => void): () => void; sendSupervisorFrame(frame: Uint8Array): Promise; onSupervisorFrame(listener: (frame: Uint8Array) => void): () => void; requestGracefulClose(signal: AbortSignal): Promise; forceTerminate(): Promise; waitForExit(deadline: number | Date): Promise; inspect(): Promise; release(): Promise; } type NodePhase = "new" | "starting" | "ready" | "failed" | "released"; /** * 受管 RPC 节点。平台适配器返回的树句柄和标准流在此处成对保存,监督器 * 只能通过本类型的高层命令与资源观察接口访问它们。 */ export class ManagedRpcNode implements ManagedRpcNodeLike { readonly process_binding = "managed" as const; private readonly adapter: ProcessTreeAdapter & { readonly launch: NonNullable; }; private readonly launchSpec: ManagedRpcNodeLaunchOptions; private readonly startupBlocked: Error | undefined; private readonly rpcOptions: Readonly> | undefined; private readonly bridgeFactory: ManagedRpcBridgeFactory; private phase: NodePhase = "new"; private binding: { readonly tree: ProcessTreeHandle; readonly transport?: ManagedProcessTransport; } | undefined; private bridge: ManagedRpcBridge | undefined; private startPromise: Promise | undefined; private readonly bindingSettled: Promise; private resolveBindingSettled!: () => void; private bindingSettlementRecorded = false; private gracefulCloseRequested = false; private forceTerminationRequested = false; private releaseRequested = false; private releasePromise: Promise | undefined; private unsubscribeBridgeEvent: (() => void) | undefined; private unsubscribeBridgeFault: (() => void) | undefined; private unsubscribeBridgeSupervisorFrame: (() => void) | undefined; constructor(options: ManagedRpcNodeOptions) { if (!isManagedProcessTreeAdapter(options.processTreeAdapter, options.processTreeAdapter.platform)) { throw new TypeError("受管 RPC 节点需要支持 launch() 的进程树适配器"); } if (!isLaunchSpec(options.launch)) throw new TypeError("受管 RPC 节点启动说明无效"); this.adapter = options.processTreeAdapter as typeof this.adapter; this.launchSpec = Object.freeze({ command: options.launch.command, ...(options.launch.args === undefined ? {} : { args: Object.freeze([...options.launch.args]) }), ...(options.launch.cwd === undefined ? {} : { cwd: options.launch.cwd }), ...(options.launch.env === undefined ? {} : { env: Object.freeze({ ...options.launch.env }) }), }); this.startupBlocked = options.startupBlocked; this.rpcOptions = options.rpcOptions === undefined ? undefined : copyRpcOptions(options.rpcOptions); this.bridgeFactory = options.bridgeFactory ?? ((transport, bridgeOptions) => new ManagedRpcBridgeClient(transport, bridgeOptions)); this.bindingSettled = new Promise((resolve) => { this.resolveBindingSettled = resolve; }); } start(signal?: AbortSignal, context?: ManagedRpcNodeStartContext): Promise { this.startPromise ??= this.runStart(signal, context); return this.startPromise; } async prompt(message: string): Promise { return this.requireBridge().prompt(message); } async steer(message: string): Promise { return this.requireBridge().steer(message); } async abort(): Promise { return this.requireBridge().abort(); } async getState(): Promise { return this.requireBridge().getState(); } onEvent(listener: (event: unknown) => void): () => void { this.eventListeners.add(listener); return () => this.eventListeners.delete(listener); } onTransportFault(listener: (fault: ManagedRpcTransportFault) => void): () => void { this.transportListeners.add(listener); return () => this.transportListeners.delete(listener); } async sendSupervisorFrame(frame: Uint8Array): Promise { await this.requireSupervisorBridge().sendSupervisorFrame(cloneBytes(frame)); } onSupervisorFrame(listener: (frame: Uint8Array) => void): () => void { this.supervisorFrameListeners.add(listener); return () => this.supervisorFrameListeners.delete(listener); } async requestGracefulClose(signal: AbortSignal): Promise { this.gracefulCloseRequested = true; await this.waitForBindingSettlement(); const bridge = this.bridge; const tree = this.binding?.tree; const operations: Promise[] = []; if (bridge !== undefined) operations.push(bridge.requestClose(signal)); if (tree !== undefined) operations.push(this.adapter.requestGracefulClose(tree, signal)); await settleAll(operations); } async forceTerminate(): Promise { this.forceTerminationRequested = true; await this.waitForBindingSettlement(); const tree = this.binding?.tree; if (tree !== undefined) await this.adapter.forceTerminate(tree); } async waitForExit(deadline: number | Date): Promise { const tree = this.binding?.tree; if (tree !== undefined) return this.adapter.waitForExit(tree, deadline); return this.phase === "starting" ? { state: "unknown" } : { state: "exited" }; } async inspect(): Promise { const tree = this.binding?.tree; if (tree !== undefined) return this.adapter.inspect(tree); return this.phase === "starting" ? { state: "unknown" } : { state: "released" }; } release(): Promise { this.releaseRequested = true; if (this.phase === "released") return Promise.resolve(); if (this.releasePromise !== undefined) return this.releasePromise; const attempt = this.runRelease(); this.releasePromise = attempt; const clearFailedAttempt = (): void => { if (this.releasePromise === attempt && this.phase !== "released") { this.releasePromise = undefined; } }; void attempt.then(() => {}, clearFailedAttempt); return attempt; } private async runRelease(): Promise { await this.waitForBindingSettlement(); if (this.phase === "released") return; const startWasInFlight = this.phase === "starting"; const tree = this.binding?.tree; const operations: Promise[] = []; if (this.bridge?.release !== undefined) operations.push(this.bridge.release()); if (tree !== undefined) operations.push(this.adapter.release(tree)); await settleAll(operations); if (startWasInFlight && this.startPromise !== undefined) { await this.startPromise.then(() => {}, () => {}); if (this.bridge?.release !== undefined) await this.bridge.release(); } this.unsubscribeBridgeEvent?.(); this.unsubscribeBridgeFault?.(); this.unsubscribeBridgeSupervisorFrame?.(); this.unsubscribeBridgeEvent = undefined; this.unsubscribeBridgeFault = undefined; this.unsubscribeBridgeSupervisorFrame = undefined; this.phase = "released"; } private async runStart(signal?: AbortSignal, context?: ManagedRpcNodeStartContext): Promise { if (this.phase !== "new") throw new Error("受管 RPC 节点已启动"); this.phase = "starting"; try { if (this.startupBlocked !== undefined) throw this.startupBlocked; if (signal?.aborted || this.cleanupRequested()) throw abortError(); const credential = randomBytes(32).toString("base64url"); const launchSpec = withBridgeCredential(this.launchSpec, credential, context?.environment); let launched: unknown; try { launched = await this.adapter.launch(launchSpec); } catch (error: unknown) { const lateTree = treeFromLaunchError(error); if (lateTree !== undefined) this.binding = Object.freeze({ tree: lateTree }); throw error; } if (launched === null || typeof launched !== "object" || !("tree" in launched)) { throw new Error("进程树启动结果无效"); } const candidate = launched as Record; if (candidate.tree === undefined) throw new Error("进程树启动结果无效"); if (!("transport" in candidate) || !isManagedTransport(candidate.transport)) { // 即使 transport 违约,也必须保留已返回的树句柄供启动回滚。 this.binding = Object.freeze({ tree: candidate.tree }); throw new Error("进程树启动结果无效"); } const binding = Object.freeze({ tree: candidate.tree, transport: candidate.transport, }); this.binding = Object.freeze({ tree: binding.tree, transport: binding.transport }); this.recordBindingSettlement(); if (signal?.aborted || this.cleanupRequested()) throw abortError(); const bridge = this.bridgeFactory(binding.transport, Object.freeze({ credential, ...(this.rpcOptions === undefined ? {} : { rpcOptions: this.rpcOptions }), })); this.bridge = bridge; this.unsubscribeBridgeEvent = bridge.onEvent((event) => this.emitEvent(event)); this.unsubscribeBridgeFault = bridge.onTransportFault((fault) => this.emitTransportFault(fault)); this.unsubscribeBridgeSupervisorFrame = bridge.onSupervisorFrame( (frame) => this.emitSupervisorFrame(frame), ); if (signal?.aborted || this.cleanupRequested()) throw abortError(); await bridge.start(signal, context); if (signal?.aborted || this.cleanupRequested() || this.hasReleased()) throw abortError(); this.phase = "ready"; } catch (error) { this.recordBindingSettlement(); if (!this.hasReleased()) this.phase = "failed"; throw error; } } private cleanupRequested(): boolean { return this.gracefulCloseRequested || this.forceTerminationRequested || this.releaseRequested; } private hasReleased(): boolean { return this.phase === "released"; } private recordBindingSettlement(): void { if (this.bindingSettlementRecorded) return; this.bindingSettlementRecorded = true; this.resolveBindingSettled(); } private async waitForBindingSettlement(): Promise { if (this.phase === "starting" && !this.bindingSettlementRecorded) { await this.bindingSettled; } } private requireBridge(): ManagedRpcBridge { if (this.phase !== "ready" || this.bridge === undefined) { throw new Error("受管 RPC 节点尚未就绪"); } return this.bridge; } private requireSupervisorBridge(): ManagedRpcBridge { if ((this.phase !== "starting" && this.phase !== "ready") || this.bridge === undefined) { throw new Error("受管 RPC 节点尚未就绪"); } return this.bridge; } private emitEvent(event: unknown): void { // 事件已经由桥接协议限制为高层对象;节点不解释 Pi JSONL。 for (const listener of this.eventListeners) { try { listener(event); } catch { // 观察者异常不能影响桥接传输。 } } } private emitTransportFault(fault: ManagedRpcTransportFault): void { for (const listener of this.transportListeners) { try { listener(fault); } catch { // 故障观察者异常不能改变资源绑定。 } } } private emitSupervisorFrame(frame: Uint8Array): void { const copy = cloneBytes(frame); for (const listener of this.supervisorFrameListeners) { try { listener(cloneBytes(copy)); } catch { // 监督帧观察者异常不能影响桥接读取者。 } } } private readonly eventListeners = new Set<(event: unknown) => void>(); private readonly transportListeners = new Set<(fault: ManagedRpcTransportFault) => void>(); private readonly supervisorFrameListeners = new Set<(frame: Uint8Array) => void>(); } function isLaunchSpec(value: unknown): value is ManagedRpcNodeLaunchOptions { if (typeof value !== "object" || value === null) return false; const candidate = value as Record; if (typeof candidate.command !== "string" || candidate.command.length === 0) return false; if (candidate.args !== undefined && ( !Array.isArray(candidate.args) || candidate.args.some((item) => typeof item !== "string") )) return false; return candidate.cwd === undefined || typeof candidate.cwd === "string"; } function isManagedTransport(value: unknown): value is ManagedProcessTransport { if (typeof value !== "object" || value === null) return false; const candidate = value as Record; return typeof candidate.stdin === "object" && candidate.stdin !== null && typeof candidate.stdout === "object" && candidate.stdout !== null && typeof candidate.stderr === "object" && candidate.stderr !== null; } function treeFromLaunchError(error: unknown): ProcessTreeHandle | undefined { if (typeof error !== "object" || error === null || !("tree" in error)) return undefined; const tree = (error as { readonly tree?: unknown }).tree; return tree === undefined ? undefined : tree; } function copyRpcOptions(value: Readonly>): Readonly> { return Object.freeze({ ...value, ...(Array.isArray(value.args) ? { args: Object.freeze([...value.args]) } : {}), }); } function cloneBytes(value: Uint8Array): Uint8Array { const copy = new Uint8Array(value.byteLength); copy.set(value); return copy; } function copySupervisorInit(value: ManagedRpcSupervisorInit): ManagedRpcSupervisorInit { return Object.freeze({ root_id: value.root_id, local_agent_id: value.local_agent_id, peer_agent_id: value.peer_agent_id, parent_agent_id: value.parent_agent_id, depth: value.depth, credential: value.credential, initial_snapshot: Object.freeze(value.initial_snapshot.map((node) => Object.freeze({ ...node }))), initial_subtree_revision: value.initial_subtree_revision, }); } async function settleAll(operations: readonly Promise[]): Promise { const results = await Promise.allSettled(operations); const failure = results.find((result): result is PromiseRejectedResult => result.status === "rejected"); if (failure !== undefined) throw failure.reason instanceof Error ? failure.reason : new Error("受管节点操作失败"); } function abortError(): Error { const error = new Error("受管 RPC 节点阶段已取消"); error.name = "AbortError"; return error; } export type ManagedRpcCommandRejectionReason = "compaction_active" | "host_busy"; /** Pi 已明确拒绝命令;与可能已经入队的传输不确定结果严格区分。 */ export class ManagedRpcCommandRejectedError extends Error { readonly reason: ManagedRpcCommandRejectionReason | undefined; constructor(reason?: ManagedRpcCommandRejectionReason) { super("受管 RPC 命令被明确拒绝"); this.name = "ManagedRpcCommandRejectedError"; this.reason = reason; } } const MAX_BRIDGE_FRAME_BYTES = MANAGED_RPC_BRIDGE_MAX_FRAME_BYTES; interface BridgeResponse { readonly protocol?: typeof MANAGED_RPC_BRIDGE_PROTOCOL; readonly kind: "response"; readonly id: number; readonly ok: boolean; readonly data?: unknown; readonly rejected?: true; readonly rejection_reason?: ManagedRpcCommandRejectionReason; readonly startup_error?: ManagedRpcStartupDiagnostic; } interface BridgeEventFrame { readonly protocol?: typeof MANAGED_RPC_BRIDGE_PROTOCOL; readonly kind: "event"; readonly event: unknown; } interface BridgeFaultFrame { readonly protocol?: typeof MANAGED_RPC_BRIDGE_PROTOCOL; readonly kind: "fault"; readonly fault: ManagedRpcTransportFault; } interface BridgeSupervisorFrame { readonly protocol?: typeof MANAGED_RPC_BRIDGE_PROTOCOL; readonly kind: "supervisor_frame"; readonly frame: string; } type BridgeFrame = BridgeResponse | BridgeEventFrame | BridgeFaultFrame | BridgeSupervisorFrame; interface PendingBridgeRequest { readonly command: string; readonly resolve: (value: unknown) => void; readonly reject: (reason: Error) => void; } /** * 受管桥接进程的父端客户端。帧以四字节大端长度前缀分隔,父端只看到高层 * 命令和安全事件;Pi JSONL 由桥接进程内部独占。 */ export interface ManagedRpcBridgeClientOptions { readonly credential?: string; /** bridge 在启动子 Pi 前应用的私有 RpcClient 配置。 */ readonly rpcOptions?: Readonly>; } export class ManagedRpcBridgeClient implements ManagedRpcBridge { private readonly stdin: Writable; private readonly stdout: Readable; private readonly stderr: Readable; private readonly pending = new Map(); private readonly eventListeners = new Set<(event: unknown) => void>(); private readonly faultListeners = new Set<(fault: ManagedRpcTransportFault) => void>(); private readonly supervisorFrameListeners = new Set<(frame: Uint8Array) => void>(); private readonly credential: string | undefined; private readonly rpcOptions: Readonly> | undefined; private readonly decoder = new LengthPrefixedFrameDecoder(MAX_BRIDGE_FRAME_BYTES); private nextRequestId = 1; private closed = false; private started = false; private startRequested = false; private writeQueue: Promise = Promise.resolve(); constructor( transport: ManagedProcessTransport, options: ManagedRpcBridgeClientOptions | string = {}, ) { this.credential = typeof options === "string" ? options : options.credential; this.rpcOptions = typeof options === "string" || options.rpcOptions === undefined ? undefined : copyRpcOptions(options.rpcOptions); this.stdin = transport.stdin; this.stdout = transport.stdout; this.stderr = transport.stderr; // bridge 会把 Pi 子进程 stderr 转发到自身 stderr;父端必须持续消费, // 否则 Windows 管道反压会冻结 bridge 的事件循环和监督 ACK。 drainStderr(transport.stderr); this.stdout.on("data", (chunk: Uint8Array | string) => this.receiveBytes( typeof chunk === "string" ? new TextEncoder().encode(chunk) : new Uint8Array(chunk), )); this.stdout.on("end", () => this.handleTransportEnd()); this.stdout.on("close", () => this.handleTransportEnd()); this.stdout.on("error", () => this.failTransport("protocol_fault")); this.stdin.on("error", () => this.failTransport("protocol_fault")); } async start(signal?: AbortSignal, context?: ManagedRpcNodeStartContext): Promise { if (this.startRequested) throw new Error("桥接进程已启动"); this.startRequested = true; await this.request( "start", context?.supervisor === undefined && this.rpcOptions === undefined ? undefined : { ...(context?.supervisor === undefined ? {} : { supervisor: copySupervisorInit(context.supervisor) }), ...(this.rpcOptions === undefined ? {} : { config: { rpc: this.rpcOptions } }), }, signal, ); this.started = true; } async prompt(message: string): Promise { await this.request("prompt", { message }); } async steer(message: string): Promise { await this.request("steer", { message }); } async abort(): Promise { await this.request("abort", undefined); } async getState(): Promise { return this.request("get_state", undefined); } async requestClose(signal: AbortSignal): Promise { await this.request("close", undefined, signal); } onEvent(listener: (event: unknown) => void): () => void { this.eventListeners.add(listener); return () => this.eventListeners.delete(listener); } onTransportFault(listener: (fault: ManagedRpcTransportFault) => void): () => void { this.faultListeners.add(listener); return () => this.faultListeners.delete(listener); } async sendSupervisorFrame(frame: Uint8Array): Promise { if (this.closed) throw new Error("桥接传输不可用"); if ( !(frame instanceof Uint8Array) || frame.byteLength === 0 || frame.byteLength > MANAGED_RPC_SUPERVISOR_MAX_ENCODED_FRAME_BYTES ) { throw new Error("监督帧无效"); } const chunks: Uint8Array[] = []; for (let offset = 0; offset < frame.byteLength; offset += MANAGED_RPC_SUPERVISOR_MAX_FRAME_BYTES) { const slice = frame.subarray(offset, Math.min(frame.byteLength, offset + MANAGED_RPC_SUPERVISOR_MAX_FRAME_BYTES)); chunks.push(encodeBridgeFrame({ protocol: MANAGED_RPC_BRIDGE_PROTOCOL, kind: "supervisor_frame", frame: Buffer.from(slice).toString("base64url"), })); } await this.enqueueWrites(chunks); } onSupervisorFrame(listener: (frame: Uint8Array) => void): () => void { this.supervisorFrameListeners.add(listener); return () => this.supervisorFrameListeners.delete(listener); } async release(): Promise { this.closed = true; for (const pending of this.pending.values()) pending.reject(new Error("桥接传输已释放")); this.pending.clear(); if (!this.stdin.destroyed) this.stdin.destroy(); if (!this.stdout.destroyed) this.stdout.destroy(); if (!this.stderr.destroyed) this.stderr.destroy(); } private request(command: string, payload: unknown, signal?: AbortSignal): Promise { if (this.closed) return Promise.reject(new Error("桥接传输不可用")); // start 帧已经入队时,close 必须能够跟在同一写入顺序域内送达 bridge; // ManagedRpcNode 可在 start 响应前进入终止流程。 if (!this.started && command !== "start" && !(command === "close" && this.startRequested)) { return Promise.reject(new Error("桥接进程尚未启动")); } const id = this.nextRequestId; this.nextRequestId += 1; if (!isBridgeCommandName(command)) return Promise.reject(new Error("桥接命令无效")); const requestPayload = command === "start" && this.credential !== undefined ? { credential: this.credential, ...(isRecord(payload) ? payload : {}) } : payload; const request = Object.freeze({ protocol: MANAGED_RPC_BRIDGE_PROTOCOL, kind: "command", id, command, ...(requestPayload === undefined ? {} : { payload: requestPayload }), }); let frames: readonly Uint8Array[]; try { frames = encodeBridgeCommandFrames(request); } catch (error) { return Promise.reject(error instanceof Error ? error : new Error("桥接命令负载无效")); } return new Promise((resolve, reject) => { const onAbort = (): void => { this.pending.delete(id); reject(abortError()); }; if (signal?.aborted) { onAbort(); return; } signal?.addEventListener("abort", onAbort, { once: true }); this.pending.set(id, { command, resolve: (value) => { signal?.removeEventListener("abort", onAbort); resolve(value); }, reject: (error) => { signal?.removeEventListener("abort", onAbort); reject(error); }, }); void this.enqueueWrites(frames).catch((error: unknown) => { this.pending.delete(id); reject(error instanceof Error ? error : new Error("桥接写入失败")); }); }); } private receiveBytes(bytes: Uint8Array): void { if (this.closed || bytes.byteLength === 0) return; try { this.decoder.push(bytes, (frameBytes) => { let parsed: unknown; try { parsed = JSON.parse(new TextDecoder("utf-8", { fatal: true }).decode(frameBytes)); } catch { this.failTransport("protocol_fault"); return false; } this.receiveFrame(parsed); return !this.closed; }); } catch { this.failTransport("protocol_fault"); } } private receiveFrame(value: unknown): void { if (!isRecord(value) || value.protocol !== MANAGED_RPC_BRIDGE_PROTOCOL || typeof value.kind !== "string") { this.failTransport("protocol_fault"); return; } if (value.kind === "response") { const hasStartupError = Object.hasOwn(value, "startup_error"); const startupError = hasStartupError ? parseManagedRpcStartupDiagnostic(value.startup_error) : undefined; if ( !hasOnlyKeys(value, [ "protocol", "kind", "id", "ok", "data", "rejected", "rejection_reason", "startup_error", ]) || !Number.isSafeInteger(value.id) || (value.id as number) <= 0 || typeof value.ok !== "boolean" || (value.ok === false && Object.hasOwn(value, "data")) || (Object.hasOwn(value, "rejected") && value.rejected !== true) || ( Object.hasOwn(value, "rejection_reason") && value.rejection_reason !== "compaction_active" && value.rejection_reason !== "host_busy" ) || (Object.hasOwn(value, "rejection_reason") && value.rejected !== true) || (hasStartupError && startupError === undefined) || (hasStartupError && (value.ok === true || value.rejected === true)) || (value.ok === true && ( Object.hasOwn(value, "rejected") || Object.hasOwn(value, "rejection_reason") || hasStartupError )) ) { this.failTransport("protocol_fault"); return; } const responseId = value.id as number; const pending = this.pending.get(responseId); if (pending === undefined) return; if (startupError !== undefined && pending.command !== "start") { this.failTransport("protocol_fault"); return; } this.pending.delete(responseId); if (value.ok) pending.resolve(value.data); else if (value.rejected === true) { pending.reject(new ManagedRpcCommandRejectedError( value.rejection_reason as ManagedRpcCommandRejectionReason | undefined, )); } else if (startupError !== undefined) { pending.reject(new ManagedRpcStartupError(startupError)); } else { pending.reject(new Error("桥接命令失败")); } return; } if (value.kind === "event") { if (!hasOnlyKeys(value, ["protocol", "kind", "event"]) || !isSafeBridgeEvent(value.event)) { this.failTransport("protocol_fault"); return; } for (const listener of this.eventListeners) { try { listener(value.event); } catch { // 观察者属于上层业务;其异常不能破坏唯一桥接读者。 } } return; } if (value.kind === "supervisor_frame") { if (!hasOnlyKeys(value, ["protocol", "kind", "frame"]) || typeof value.frame !== "string") { this.failTransport("protocol_fault"); return; } const frame = decodeBase64Bytes(value.frame); if ( frame === undefined || frame.byteLength === 0 || frame.byteLength > MANAGED_RPC_SUPERVISOR_MAX_FRAME_BYTES ) { this.failTransport("protocol_fault"); return; } for (const listener of this.supervisorFrameListeners) { try { listener(cloneBytes(frame)); } catch { // 观察者异常不能破坏唯一传输读取者。 } } return; } if ( value.kind === "fault" && hasOnlyKeys(value, ["protocol", "kind", "fault"]) && (value.fault === "eof" || value.fault === "protocol_fault" || value.fault === "process_exit") ) { this.failTransport(value.fault as ManagedRpcTransportFault); return; } this.failTransport("protocol_fault"); } private failTransport(fault: ManagedRpcTransportFault): void { if (this.closed) return; this.closed = true; this.decoder.reset(); for (const pending of this.pending.values()) pending.reject(new Error("桥接传输故障")); this.pending.clear(); for (const listener of this.faultListeners) { try { listener(fault); } catch { // 故障观察者异常不能再次进入桥接故障路径。 } } } private handleTransportEnd(): void { this.failTransport(this.decoder.hasPendingBytes() ? "protocol_fault" : "eof"); } private enqueueWrites(frames: readonly Uint8Array[]): Promise { const operation = this.writeQueue.catch(() => {}).then(async () => { for (const frame of frames) await writeChunk(this.stdin, frame); }); this.writeQueue = operation.catch(() => {}); return operation; } } function encodeBridgeCommandFrames(request: Readonly>): readonly Uint8Array[] { try { return Object.freeze([encodeBridgeFrame(request)]); } catch (error) { if (!(error instanceof Error) || error.message !== "桥接帧超限") throw error; } if (!Object.hasOwn(request, "payload")) throw new Error("桥接命令负载无效"); if ( !Number.isSafeInteger(request.id) || (request.id as number) <= 0 || typeof request.command !== "string" || !isBridgeCommandName(request.command) ) throw new Error("桥接命令无效"); const payload = encodeBridgeJson(request.payload); if (payload.byteLength === 0 || payload.byteLength > MANAGED_RPC_BRIDGE_MAX_COMMAND_PAYLOAD_BYTES) { throw new Error("桥接命令负载超限"); } const chunkCount = Math.ceil(payload.byteLength / MANAGED_RPC_BRIDGE_COMMAND_CHUNK_BYTES); const frames: Uint8Array[] = []; for (let offset = 0, index = 0; offset < payload.byteLength; offset += MANAGED_RPC_BRIDGE_COMMAND_CHUNK_BYTES, index += 1) { const chunk = payload.subarray(offset, Math.min(payload.byteLength, offset + MANAGED_RPC_BRIDGE_COMMAND_CHUNK_BYTES)); frames.push(encodeBridgeFrame({ protocol: MANAGED_RPC_BRIDGE_PROTOCOL, kind: "command_chunk", id: request.id, command: request.command, chunk_index: index, chunk_count: chunkCount, chunk: Buffer.from(chunk).toString("base64url"), })); } return Object.freeze(frames); } function encodeBridgeFrame(value: unknown): Uint8Array { const body = encodeBridgeJson(value); if (body.byteLength > MAX_BRIDGE_FRAME_BYTES) throw new Error("桥接帧超限"); const frame = new Uint8Array(body.byteLength + 4); new DataView(frame.buffer).setUint32(0, body.byteLength, false); frame.set(body, 4); return frame; } function encodeBridgeJson(value: unknown): Uint8Array { let text: string | undefined; try { text = JSON.stringify(value); } catch { throw new Error("桥接帧无效"); } if (typeof text !== "string") throw new Error("桥接帧无效"); return new TextEncoder().encode(text); } function drainStderr(stream: Readable): void { stream.on("data", () => {}); stream.on("error", () => {}); stream.resume(); } function writeChunk(stream: Writable, bytes: Uint8Array): Promise { return new Promise((resolve, reject) => { try { stream.write(bytes, (error?: Error | null) => error === null || error === undefined ? resolve() : reject(error)); } catch (error) { reject(error instanceof Error ? error : new Error("桥接写入失败")); } }); } function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } function hasOnlyKeys(value: Record, allowed: readonly string[]): boolean { return Object.keys(value).every((key) => allowed.includes(key)); } const BRIDGE_COMMAND_NAMES = new Set([ "start", "prompt", "steer", "abort", "get_state", "close", ]); function isBridgeCommandName(value: string): boolean { return BRIDGE_COMMAND_NAMES.has(value); } function isSafeBridgeEvent(value: unknown): boolean { if (!isRecord(value) || typeof value.type !== "string") return false; switch (value.type) { case "agent_start": case "agent_settled": return Object.keys(value).length === 1; case "compaction_start": return (value.reason === "manual" || value.reason === "threshold" || value.reason === "overflow") && Object.keys(value).every((key) => key === "type" || key === "reason"); case "compaction_end": return (value.reason === "manual" || value.reason === "threshold" || value.reason === "overflow") && typeof value.aborted === "boolean" && typeof value.willRetry === "boolean" && typeof value.failed === "boolean" && Object.keys(value).every((key) => [ "type", "reason", "aborted", "willRetry", "failed", ].includes(key)); case "queue_update": return Number.isSafeInteger(value.pendingMessageCount) && (value.pendingMessageCount as number) >= 0 && Object.keys(value).every((key) => key === "type" || key === "pendingMessageCount"); case "message": return isSafeActivityMessageEvent(value); case "model_call_failure": return isSafeModelCallFailureEvent(value); case "tool_execution_start": case "tool_execution_end": return isSafeActivityToolEvent(value); case "extension_error": return Object.keys(value).length === 1; default: return false; } } function isSafeActivityMessageEvent(value: Record): boolean { if ( !Array.isArray(value.content) || value.content.length === 0 || value.content.length > 64 || !Object.keys(value).every((key) => key === "type" || key === "content") ) return false; for (const item of value.content) { if (!isRecord(item) || typeof item.type !== "string") return false; if (item.type === "text") { // assistant 正文聚合无字节上限;帧边界由长度前缀帧保证。 if (!isActivityTextShape(item.text)) return false; continue; } if (item.type === "thinking") { if (!isActivityTextShape(item.thinking)) return false; continue; } return false; } return true; } function isSafeActivityToolEvent(value: Record): boolean { if ( typeof value.toolCallId !== "string" || value.toolCallId.length === 0 || value.toolCallId.length > 256 || typeof value.toolName !== "string" || value.toolName.length === 0 || value.toolName.length > 256 ) return false; const allowed = value.type === "tool_execution_start" ? ["type", "toolCallId", "toolName", "origin", "executionGeneration"] : ["type", "toolCallId", "toolName", "origin", "executionGeneration", "isError"]; if (!Object.keys(value).every((key) => allowed.includes(key))) return false; if ( value.executionGeneration !== undefined && !isValidToolExecutionGeneration(value.executionGeneration) ) return false; if (!isSafeToolOrigin(value.origin)) return false; return value.type === "tool_execution_start" || typeof value.isError === "boolean"; } /** * 模型调用失败事件按固定字段集合校验:桥接帧闭集与产生端归一化保持同一 * 形状,避免同一事件在发布侧合法、在接收侧变成通道故障。桥接 RPC 副本的 * 失败事实始终带 provider/model;无身份的压缩自身失败只走子代理产生端的 * 监督活动路径,不经过该闭集。 */ function isSafeModelCallFailureEvent(value: Record): boolean { if (!Object.keys(value).every((key) => [ "type", "failure", "message", "provider", "model", ].includes(key))) return false; if (!isSafeModelCallFailureReason(value.failure)) return false; return isActivityTextShape(value.message) && value.message.length > 0 && isBoundedIdentityShape(value.provider) && isBoundedIdentityShape(value.model); } /** provider 与 model 身份是短引用:非空且不超过工具身份同款上限。 */ function isBoundedIdentityShape(value: unknown): value is string { return typeof value === "string" && value.length > 0 && value.length <= 256; } /** * assistant 消息正文按类型粗校验:正文聚合无字节上限,精确帧预算由桥接 * 长度前缀帧与监督通道分块共同保证。 */ function isActivityTextShape(value: unknown): value is string { return typeof value === "string"; } function decodeBase64Bytes(value: string): Uint8Array | undefined { if (value.length === 0 || !/^[A-Za-z0-9_-]+={0,2}$/.test(value)) return undefined; const padding = value.endsWith("==") ? 2 : value.endsWith("=") ? 1 : 0; if (padding > 0 && value.length % 4 !== 0) return undefined; if ((value.length - padding) % 4 === 1) return undefined; try { const normalized = value.replace(/-/g, "+").replace(/_/g, "/"); const padded = padding > 0 ? normalized : normalized + "=".repeat((4 - (normalized.length % 4)) % 4); const bytes = Buffer.from(padded, "base64"); if (bytes.toString("base64url") !== value.replace(/=+$/, "")) return undefined; return new Uint8Array(bytes); } catch { return undefined; } } function withBridgeCredential( spec: ManagedRpcNodeLaunchOptions, credential: string, runtimeEnvironment?: Readonly>, ): ManagedRpcNodeLaunchOptions { return Object.freeze({ ...spec, env: Object.freeze({ ...process.env, ...(spec.env ?? {}), ...(runtimeEnvironment ?? {}), [MANAGED_RPC_BRIDGE_CREDENTIAL_ENV]: credential, }), }); } export interface FakeManagedRpcNodeOptions { readonly onOperation?: (operation: string) => void; readonly state?: unknown; } /** 供监督器与旅程测试使用的确定性受管节点替身。 */ export class FakeManagedRpcNode implements ManagedRpcNodeLike { readonly process_binding = "managed" as const; private readonly options: FakeManagedRpcNodeOptions; private readonly eventListeners = new Set<(event: unknown) => void>(); private readonly faultListeners = new Set<(fault: ManagedRpcTransportFault) => void>(); private readonly supervisorFrameListeners = new Set<(frame: Uint8Array) => void>(); private readonly operationLog: string[] = []; private phase: "new" | "ready" | "released" = "new"; constructor(options: FakeManagedRpcNodeOptions = {}) { this.options = options; } async start(_signal?: AbortSignal, _context?: ManagedRpcNodeStartContext): Promise { this.record("start"); this.phase = "ready"; } async prompt(): Promise { this.record("prompt"); } async steer(): Promise { this.record("steer"); } async abort(): Promise { this.record("abort"); } async getState(): Promise { this.record("get_state"); return this.options.state ?? { isStreaming: false, isCompacting: false, pendingMessageCount: 0, }; } onEvent(listener: (event: unknown) => void): () => void { this.eventListeners.add(listener); return () => this.eventListeners.delete(listener); } onTransportFault(listener: (fault: ManagedRpcTransportFault) => void): () => void { this.faultListeners.add(listener); return () => this.faultListeners.delete(listener); } async sendSupervisorFrame(_frame: Uint8Array): Promise { this.record("supervisor_frame"); } onSupervisorFrame(listener: (frame: Uint8Array) => void): () => void { this.supervisorFrameListeners.add(listener); return () => this.supervisorFrameListeners.delete(listener); } async requestGracefulClose(): Promise { this.record("graceful_close"); } async forceTerminate(): Promise { this.record("force_terminate"); } async waitForExit(): Promise { this.record("wait_for_exit"); return { state: "exited" }; } async inspect(): Promise { this.record("inspect"); return { state: "released" }; } async release(): Promise { this.record("release"); this.phase = "released"; } emitEvent(event: unknown): void { for (const listener of this.eventListeners) listener(event); } emitTransportFault(fault: ManagedRpcTransportFault): void { for (const listener of this.faultListeners) listener(fault); } emitSupervisorFrame(frame: Uint8Array): void { for (const listener of this.supervisorFrameListeners) listener(cloneBytes(frame)); } operations(): readonly string[] { return Object.freeze([...this.operationLog]); } private record(operation: string): void { this.operationLog.push(operation); this.options.onOperation?.(operation); } }