import type { Readable, Writable } from "node:stream"; import type { CanonicalAgentActivityEntry } from "./canonical-activity.ts"; import type { SafeAgentActivityDisplayEvent } from "./rpc-bridge-event.ts"; import { SupervisorChannel, SupervisorFrameDecoder, type SupervisorActivityDelivery, type SupervisorChannelOptions, type SupervisorCapabilityManifest, type SupervisorControlRequest, type SupervisorControlResponse, type SupervisorDisplayDelivery, type SupervisorFrame, type SupervisorEvent, type SupervisorReceiveResult, type SupervisorReply, type SupervisorSnapshot, type SupervisorChannelPublicState, } from "./supervisor-channel.ts"; import type { RpcSupervisorChannel, RpcSupervisorChannelCloseState, RpcSupervisorChannelFault, } from "./rpc-supervisor.ts"; import { createDeferred, notifySupervisorListeners, raceSupervisorAbort, supervisorAbortError, waitForSupervisorSignal, type Deferred, } from "./supervisor-channel-async.ts"; export interface SupervisorByteTransport { readonly stdin: Writable; readonly stdout: Readable; } export interface StreamSupervisorChannelOptions extends SupervisorChannelOptions { readonly transport: SupervisorByteTransport; /** child 端握手后发送的首个完整安全快照。 */ readonly initialSnapshot?: readonly unknown[]; readonly initialSubtreeRevision?: number; /** 父端收到已经通过协议校验的生命周期事实。 */ readonly onEvent?: (event: SupervisorEvent) => void; readonly onSnapshot?: (snapshot: SupervisorSnapshot) => void; /** parent 在普通 ready 后收到的一次性内部 capability manifest。 */ readonly onCapability?: (capability: SupervisorCapabilityManifest) => void; /** child 收到 close 后执行后代优先清理;只有 true 才关闭本地字节流。 */ readonly onCloseRequested?: () => boolean | Promise; /** child 等待父端同步返回扩展消息提交结果的内部期限。 */ readonly replyDispatchTimeoutMs?: number; } /** * 将 SupervisorChannel 协议状态机绑定到一条可靠字节流。该类只负责传输和 * 期限,不替调用方裁决生命周期;协议故障通过稳定的 eof/protocol_fault 上报。 */ export class StreamSupervisorChannel implements RpcSupervisorChannel { private readonly transport: SupervisorByteTransport; private readonly protocol: SupervisorChannel; private readonly decoder: SupervisorFrameDecoder; private readonly initialSnapshot: readonly unknown[]; private readonly initialSubtreeRevision: number | undefined; private readonly onCloseRequested: (() => boolean | Promise) | undefined; private readonly replyDispatchTimeoutMs: number; private readonly ready = createDeferred(); private readonly closed = createDeferred(); private readonly replyDispatches = new Map>(); private readonly faults = new Set<(fault: RpcSupervisorChannelFault) => void>(); private readonly eventListeners = new Set<(event: SupervisorEvent) => void>(); private readonly activityListeners = new Set<(activity: SupervisorActivityDelivery) => void>(); private readonly displayListeners = new Set<(delivery: SupervisorDisplayDelivery) => void>(); private readonly snapshotListeners = new Set<(snapshot: SupervisorSnapshot) => void>(); private readonly capabilityListeners = new Set<(capability: SupervisorCapabilityManifest) => void>(); private readonly controlRequestListeners = new Set<(request: SupervisorControlRequest) => void>(); private readonly controlResponseListeners = new Set<(response: SupervisorControlResponse) => void>(); private writeQueue: Promise = Promise.resolve(); private deferredFrameSends: Array<{ readonly frame: SupervisorFrame; readonly completion: Deferred; }> | undefined; private bound = false; private snapshotSent = false; private released = false; private faultNotified = false; private endpointClosed = false; private closeHandling: Promise | undefined; constructor(options: StreamSupervisorChannelOptions) { this.transport = options.transport; this.protocol = new SupervisorChannel({ role: options.role, rootId: options.rootId, localAgentId: options.localAgentId, peerAgentId: options.peerAgentId, parentAgentId: options.parentAgentId, depth: options.depth, credential: options.credential, requestIdRegistry: options.requestIdRegistry, ...(options.limits === undefined ? {} : { limits: options.limits }), ...(options.streamIdFactory === undefined ? {} : { streamIdFactory: options.streamIdFactory }), ...(options.onReply === undefined ? {} : { onReply: options.onReply }), ...(options.resyncTimeoutMs === undefined ? {} : { resyncTimeoutMs: options.resyncTimeoutMs }), onProtocolFault: () => this.fail("protocol_fault"), }); this.decoder = new SupervisorFrameDecoder(options.limits); this.initialSnapshot = Object.freeze([...(options.initialSnapshot ?? [])]); this.initialSubtreeRevision = options.initialSubtreeRevision; this.onCloseRequested = options.onCloseRequested; this.replyDispatchTimeoutMs = validPositiveDuration(options.replyDispatchTimeoutMs) ? options.replyDispatchTimeoutMs : validPositiveDuration(options.resyncTimeoutMs) ? options.resyncTimeoutMs : 5_000; // child 可能在父端等待握手前先收到 EOF;不能让内部 ready 拒绝升级为 // 未处理拒绝,但 waitForReady 仍需观察原始失败。 void this.ready.promise.catch(() => {}); if (options.onEvent !== undefined) this.eventListeners.add(options.onEvent); if (options.onSnapshot !== undefined) this.snapshotListeners.add(options.onSnapshot); if (options.onCapability !== undefined) this.capabilityListeners.add(options.onCapability); this.transport.stdout.on("data", (chunk: Uint8Array | string) => { try { const bytes = typeof chunk === "string" ? new TextEncoder().encode(chunk) : new Uint8Array(chunk); for (const frame of this.decoder.push(bytes)) this.receive(frame); } catch { this.protocol.markProtocolFault(); this.fail("protocol_fault"); } }); this.transport.stdout.on("end", () => this.onEof()); this.transport.stdout.on("close", () => this.onEof()); this.transport.stdout.on("error", () => this.fail("protocol_fault")); this.transport.stdin.on("error", () => this.fail("protocol_fault")); } async bind(signal: AbortSignal): Promise { if (this.bound) return; this.bound = true; if (signal.aborted) throw supervisorAbortError(); if (this.protocol.role === "child") { await this.send(this.protocol.startHandshake()); } } async waitForReady(signal: AbortSignal): Promise { if (this.isReady()) return; if (signal.aborted) throw supervisorAbortError(); await raceSupervisorAbort(this.ready.promise, signal); } isReady(): boolean { return this.publicState().state === "ready"; } /** parent reload 建立 reset snapshot 边界;窗口内 activity/display 可安全丢弃。 */ async requestSnapshot(): Promise { await this.send(this.protocol.requestSnapshot()); } async publishReply( reply: SupervisorReply, signal?: AbortSignal, ): Promise { if (signal?.aborted === true) throw supervisorAbortError(); const frame = this.protocol.publishReply(reply); const requestId = frame.request_id; if (requestId === undefined) throw new Error("message_delivery_failed"); const waiter = createDeferred(); void waiter.promise.catch(() => {}); this.replyDispatches.set(requestId, waiter); try { await this.send(frame); const accepted = await this.waitForReplyDispatch(waiter, signal); if (!accepted) throw new Error("message_delivery_failed"); } finally { if (this.replyDispatches.get(requestId) === waiter) this.replyDispatches.delete(requestId); } } /** child 在普通 ready 后发布一次内部 capability manifest。 */ async publishCapability(capability: SupervisorCapabilityManifest): Promise { await this.send(this.protocol.publishCapability(capability)); } /** child 端发布已经由 SupervisorChannel 校验的生命周期事实。 */ async publishEvent( event: Omit & { readonly agent_id?: string }, ): Promise { await this.send(this.protocol.publishEvent(event)); } /** * child 端发布规范活动条目;超过单帧预算的条目由协议层分块为多帧逐帧 * 发送。发布端拒绝时静默返回,会话不受影响。 */ async publishActivity(input: { readonly agent_id?: string; readonly entry: CanonicalAgentActivityEntry; }): Promise { const frames = this.protocol.publishActivity(input); for (const frame of frames) await this.send(frame); } /** * child 端发布实时显示事件;fire-and-forget,发布端拒绝时静默返回,会话 * 不受影响。 */ async publishDisplayActivity(input: { readonly agent_id?: string; readonly event: SafeAgentActivityDisplayEvent; }): Promise { const frames = this.protocol.publishDisplayActivity(input); for (const frame of frames) await this.send(frame); } /** child 端发布新的完整子树;修订和正文边界仍由协议状态机校验。 */ async publishSnapshot(nodes: readonly unknown[], subtreeRevision: number): Promise { await this.send(this.protocol.publishSnapshot(nodes, subtreeRevision)); } async publishControlRequest(request: SupervisorControlRequest): Promise { await this.send(this.protocol.publishControlRequest(request)); } async publishControlResponse(response: SupervisorControlResponse): Promise { await this.send(this.protocol.publishControlResponse(response)); } establishTerminationBarrier(): void { this.protocol.establishTerminationBarrier(); } async requestClose(signal: AbortSignal): Promise { if (signal.aborted) throw supervisorAbortError(); try { await this.send(this.protocol.createCloseFrame()); } catch { // 关闭帧无法送达时,资源观察仍由受管节点/平台适配器裁决。 } } async waitForClose(deadline: number | Date): Promise { if (this.released) return "unknown"; const end = deadline instanceof Date ? deadline.getTime() : deadline; const remaining = Math.max(0, end - Date.now()); if (!this.closed.settled() && remaining > 0) { await waitForSupervisorSignal(this.closed.promise, remaining); } return this.endpointClosed ? "released" : "present"; } async release(): Promise { if (this.released) return; this.released = true; this.endpointClosed = true; for (const waiter of this.replyDispatches.values()) waiter.reject(new Error("监督通道已关闭")); this.replyDispatches.clear(); try { if (!this.transport.stdin.destroyed) this.transport.stdin.destroy(); if (!this.transport.stdout.destroyed) this.transport.stdout.destroy(); } finally { this.closed.resolve(); } } onFault(listener: (fault: RpcSupervisorChannelFault) => void): () => void { this.faults.add(listener); return () => this.faults.delete(listener); } onEvent(listener: (event: SupervisorEvent) => void): () => void { this.eventListeners.add(listener); return () => this.eventListeners.delete(listener); } onActivity(listener: (activity: SupervisorActivityDelivery) => void): () => void { this.activityListeners.add(listener); return () => this.activityListeners.delete(listener); } onDisplay(listener: (delivery: SupervisorDisplayDelivery) => void): () => void { this.displayListeners.add(listener); return () => this.displayListeners.delete(listener); } onSnapshot(listener: (snapshot: SupervisorSnapshot) => void): () => void { this.snapshotListeners.add(listener); return () => this.snapshotListeners.delete(listener); } /** 获取 parent 已缓存的一次性 capability;该信息不属于公开状态。 */ getCapability(): SupervisorCapabilityManifest | undefined { return this.protocol.getCapability(); } /** 注册 capability 观察者;若已缓存则同步交付安全副本。 */ onCapability(listener: (capability: SupervisorCapabilityManifest) => void): () => void { this.capabilityListeners.add(listener); const capability = this.protocol.getCapability(); if (capability !== undefined) notifySupervisorListeners(new Set([listener]), capability); return () => this.capabilityListeners.delete(listener); } onControlRequest(listener: (request: SupervisorControlRequest) => void): () => void { this.controlRequestListeners.add(listener); return () => this.controlRequestListeners.delete(listener); } onControlResponse(listener: (response: SupervisorControlResponse) => void): () => void { this.controlResponseListeners.add(listener); return () => this.controlResponseListeners.delete(listener); } /** 路由相关性或 operation_id 复用违约时固定为监督协议故障。 */ failProtocol(): void { this.protocol.markProtocolFault(); this.fail("protocol_fault"); } /** 只读协议状态用于测试/诊断,不暴露凭据或端点。 */ getPublicState(): SupervisorChannelPublicState { return this.protocol.getPublicState(); } private receive(frame: SupervisorFrame): void { if (this.deferredFrameSends !== undefined) { this.protocol.markProtocolFault(); this.fail("protocol_fault"); return; } this.deferredFrameSends = []; const result = this.protocol.receive(frame); if (result.kind === "eof" || result.kind === "protocol_fault") { this.rejectDeferredFrameSends(); this.fail(result.kind === "eof" ? "eof" : "protocol_fault"); return; } // onReply 在 protocol.receive() 内同步执行,可能先创建下一条业务帧; // 记录边界,确保这些较早分配序号的帧先于本次协议 ACK 写出。 const framesBeforeProtocolResponses = this.deferredFrameSends?.length ?? 0; if (result.kind === "accepted" && result.event !== undefined) { for (const listener of this.eventListeners) { try { listener(result.event); } catch { // 观察者异常不能改变协议状态或破坏后续帧读取。 } } } if (result.kind === "accepted" && result.activity !== undefined) { for (const listener of this.activityListeners) { try { listener(result.activity); } catch { // 活动观察者异常不能改变协议状态。 } } } if (result.kind === "accepted" && result.display !== undefined) { for (const listener of this.displayListeners) { try { listener(result.display); } catch { // 显示观察者异常不能改变协议状态。 } } } if (result.kind === "accepted" && result.snapshot !== undefined) { for (const listener of this.snapshotListeners) { try { listener(result.snapshot); } catch { // 快照观察者异常不能回滚协议已接受的原子替换。 } } } if (result.kind === "accepted" && result.capability !== undefined) { notifySupervisorListeners(this.capabilityListeners, result.capability); } if (result.kind === "accepted" && result.control_request !== undefined) { notifySupervisorListeners(this.controlRequestListeners, result.control_request); } if (result.kind === "accepted" && result.control_response !== undefined) { const dispatch = readReplyDispatch(result.control_response); if (dispatch === undefined) { notifySupervisorListeners(this.controlResponseListeners, result.control_response); } else { this.resolveReplyDispatch(dispatch.operationId, dispatch.accepted); } } if (result.kind === "accepted" && result.close_requested === true) { this.handleCloseRequested(); } if (this.protocol.role === "child" && !this.snapshotSent && this.protocol.getPublicState().state === "awaiting_snapshot") { this.snapshotSent = true; try { void this.send(this.protocol.publishSnapshot(this.initialSnapshot, this.initialSubtreeRevision)); } catch { this.rejectDeferredFrameSends(); this.fail("protocol_fault"); return; } } const protocolSends = result.kind === "accepted" || result.kind === "duplicate" || result.kind === "gap" ? this.flushDeferredFrameSends(result.outbound, framesBeforeProtocolResponses) : this.flushDeferredFrameSends([], framesBeforeProtocolResponses); if ( (result.kind === "accepted" || result.kind === "duplicate" || result.kind === "gap") && this.isReady() ) { void Promise.all(protocolSends).then( () => this.ready.resolve(), () => this.fail("protocol_fault"), ); } if (this.isReady()) this.ready.resolve(); } private publicState(): SupervisorChannelPublicState { return this.protocol.getPublicState(); } private send(frame: SupervisorFrame): Promise { if (this.deferredFrameSends !== undefined) { const completion = createDeferred(); void completion.promise.catch(() => {}); this.deferredFrameSends.push({ frame, completion }); return completion.promise; } return this.sendDirect(frame); } private async sendDirect(frame: SupervisorFrame): Promise { const encoded = this.protocol.encode(frame); const operation = this.writeQueue.catch(() => {}).then(() => new Promise((resolve, reject) => { try { this.transport.stdin.write(encoded, (error?: Error | null) => { if (error === undefined || error === null) resolve(); else reject(error); }); } catch (error) { reject(error instanceof Error ? error : new Error("监督帧写入失败")); } })); // 保留一个已恢复的尾部,单次写入失败不能毒化后续控制帧。 this.writeQueue = operation.catch(() => {}); return operation; } private flushDeferredFrameSends( protocolFrames: readonly SupervisorFrame[], framesBeforeProtocolResponses: number, ): readonly Promise[] { const pending = this.deferredFrameSends ?? []; this.deferredFrameSends = undefined; const sendDeferred = (items: typeof pending): void => { for (const item of items) { void this.sendDirect(item.frame).then(item.completion.resolve, item.completion.reject); } }; sendDeferred(pending.slice(0, framesBeforeProtocolResponses)); const protocolSends = protocolFrames.map((frame) => { const send = this.sendDirect(frame); void send.catch(() => this.fail("protocol_fault")); return send; }); sendDeferred(pending.slice(framesBeforeProtocolResponses)); return Object.freeze(protocolSends); } private rejectDeferredFrameSends(): void { const pending = this.deferredFrameSends ?? []; this.deferredFrameSends = undefined; for (const item of pending) item.completion.reject(new Error("监督通道不可用")); } private resolveReplyDispatch(operationId: string, accepted: boolean): void { const waiter = this.replyDispatches.get(operationId); if (waiter === undefined) return; this.replyDispatches.delete(operationId); waiter.resolve(accepted); } private async waitForReplyDispatch( waiter: Deferred, signal?: AbortSignal, ): Promise { let timer: ReturnType | undefined; const timeout = new Promise((_, reject) => { timer = setTimeout(() => reject(new Error("message_delivery_failed")), this.replyDispatchTimeoutMs); timer.unref?.(); }); try { const pending = Promise.race([waiter.promise, timeout]); return signal === undefined ? await pending : await raceSupervisorAbort(pending, signal); } finally { if (timer !== undefined) clearTimeout(timer); } } private onEof(): void { if (this.released) return; if (this.endpointClosed) return; this.endpointClosed = true; try { this.decoder.finish(); } catch { this.protocol.markProtocolFault(); this.closed.resolve(); this.fail("protocol_fault"); return; } this.protocol.receiveEof(); this.closed.resolve(); this.fail("eof"); } private handleCloseRequested(): void { if (this.closeHandling !== undefined || this.onCloseRequested === undefined) return; const operation = Promise.resolve() .then(() => this.onCloseRequested?.() ?? false) .then(async (complete) => { if (complete === true) await this.release(); }) .catch(() => { // 未确认清理时保持传输存在,由父端内部期限和平台树回收继续裁决。 }); this.closeHandling = operation; } private fail(fault: RpcSupervisorChannelFault): void { if (this.released) return; if (this.faultNotified) return; this.faultNotified = true; this.closed.resolve(); for (const waiter of this.replyDispatches.values()) waiter.reject(new Error("监督通道不可用")); this.replyDispatches.clear(); if (!this.ready.settled()) this.ready.reject(new Error("监督通道不可用")); for (const listener of this.faults) { try { listener(fault); } catch { // 观察者异常不能再次破坏通道状态。 } } } } function readReplyDispatch( response: SupervisorControlResponse, ): { readonly operationId: string; readonly accepted: boolean } | undefined { if (!response.ok || typeof response.data !== "object" || response.data === null || Array.isArray(response.data)) { return undefined; } const data = response.data as Record; if ( data.kind !== "reply_dispatch" || typeof data.accepted !== "boolean" || Object.keys(data).some((key) => key !== "kind" && key !== "accepted") ) return undefined; return Object.freeze({ operationId: response.operation_id, accepted: data.accepted, }); } function validPositiveDuration(value: number | undefined): value is number { return value !== undefined && Number.isSafeInteger(value) && value > 0; }