/** * Replay Buffer — 每会话流式回放缓冲(v0.8.0 P1-1a 引入;v0.8.1 P0-3 模块化收口) * * 渲染层切走会话期间,事件管道仍照常投递(按会话隔离被渲染层准入拒绝), * 此缓冲按序保存 streamEvent + stateChange 双通道事件;用户切回会话时经 * agent:getReplayState 拉取并灌入渲染层事件总线,实现"后台会话运行内容不丢失"。 * * v0.8.1 P0-3 根治:缓冲原为 ipc/agent.ts 内部 Map —— 会话删除/彻底删除时无 * 联动清理,仅靠 50 会话 LRU 兜底,已删会话的缓冲(最多 4MB/会话)滞留内存。 * 现抽为独立模块,sessions:delete / sessions:purge 终态路径显式清除。 * * 有界设计:单会话上限 MAX_REPLAY_EVENTS 条 / MAX_REPLAY_BYTES 字节,溢出丢弃 * 最旧并置 truncated 标记;TERMINATED 后保留(供"完成后切回"看到最终内容), * 下一次 run 启动(INIT stateChange)时清空重建。 */ const MAX_REPLAY_EVENTS = 2000; const MAX_REPLAY_BYTES = 4 * 1024 * 1024; const MAX_REPLAY_SESSIONS = 50; export interface ReplayEntry { channel: 'streamEvent' | 'stateChange'; payload: unknown; ts: number; } export interface ReplayBuffer { events: ReplayEntry[]; bytes: number; truncated: boolean; /** 最近一次观察到的 runId(供渲染层回放后对齐 run 守卫) */ runId?: string; terminated: boolean; } const replayBuffers = new Map(); /** 追加一条事件到会话回放缓冲(有界:条数/字节双上限,溢出丢最旧) */ export function appendReplay(sessionId: string, entry: ReplayEntry): void { let buf = replayBuffers.get(sessionId); if (!buf) { // 会话数上限保护:超出后淘汰最早的缓冲(Map 迭代序 = 插入序) if (replayBuffers.size >= MAX_REPLAY_SESSIONS) { const oldest = replayBuffers.keys().next().value as string | undefined; if (oldest !== undefined) replayBuffers.delete(oldest); } buf = { events: [], bytes: 0, truncated: false, terminated: false }; replayBuffers.set(sessionId, buf); } // run 启动(INIT)后旧 run 缓冲作废 —— 由 stateChange 监听器在 INIT 时清空, // 此处仅在 terminated 缓冲上遇到新内容时兜底重置(正常路径 INIT 先行) if (buf.terminated) { replayBuffers.delete(sessionId); buf = { events: [], bytes: 0, truncated: false, terminated: false }; replayBuffers.set(sessionId, buf); } let size = 0; try { size = JSON.stringify(entry.payload).length; } catch { size = 256; // 序列化失败按保守值计入 } buf.events.push(entry); buf.bytes += size; while ( buf.events.length > MAX_REPLAY_EVENTS || (buf.bytes > MAX_REPLAY_BYTES && buf.events.length > 1) ) { const dropped = buf.events.shift(); buf.truncated = true; if (dropped) { try { buf.bytes -= JSON.stringify(dropped.payload).length; } catch { buf.bytes -= 256; } } } // 记录 runId(首个携带者) const payload = entry.payload as { runId?: string } | null; if (!buf.runId && payload?.runId) buf.runId = payload.runId; } /** run 启动(INIT)时清空上一 run 的缓冲 */ export function resetReplayBuffer(sessionId: string): void { replayBuffers.delete(sessionId); } /** TERMINATED 时置终态标记(保留缓冲供"完成后切回"回放最终内容) */ export function markReplayTerminated(sessionId: string): void { const buf = replayBuffers.get(sessionId); if (buf) buf.terminated = true; } /** v0.8.1 P0-3: 会话终态(删除/彻底删除)时清除缓冲,杜绝内存滞留 */ export function clearReplayBuffer(sessionId: string): void { replayBuffers.delete(sessionId); } /** 拉取会话回放状态(agent:getReplayState 消费) */ export function getReplayBufferData(sessionId: string): { runId: string | null; truncated: boolean; events: ReplayEntry[]; } { const buf = replayBuffers.get(sessionId); return { runId: buf?.runId ?? null, truncated: buf?.truncated ?? false, events: buf?.events ?? [], }; }