/** * Agent Engine Manager — 每会话独立引擎管理器(P2-10) * * 解决原"全局单引擎"的两个缺陷: * 1. 全局串行锁:原 AgentLoopEngine.currentRunPromise 使所有会话共享一把锁, * 上一会话未结束时新会话消息需排队等待(最长卡 120s 工具超时)。 * 现在每个会话持有独立引擎实例,多会话可并行运行。 * 2. adapter abort 信号互踩:原所有引擎/SubAgent 共享一个 adapter 实例, * setAbortSignal 单槽位导致并发时中断信号错乱。 * 现在创建引擎时通过 adapter 工厂为每个引擎生成独立 adapter 实例 * (adapter 是无状态的配置包装,实例化成本可忽略)。 * * 引擎生命周期: * - 按需创建(首次 sendMessage 时),事件统一转发到 manager(附加 sessionId) * - LRU 淘汰:缓存超过 30 个引擎时,淘汰最旧的非运行中引擎 * * @see electron/harness/agent-loop/engine.ts — 引擎实现 */ import { EventEmitter } from 'events'; import { AgentLoopEngine } from '../harness/agent-loop'; import type { AgentLoopConfig } from '../harness/agent-loop/types'; import type { IMetonaProviderAdapter, MetonaToolDef } from '../harness/types'; import type { ToolRegistry } from '../harness/tools/registry'; import type { PreToolHook } from '../harness/hooks/pre-tool'; import type { PostToolHook } from '../harness/hooks/post-tool'; import log from 'electron-log'; /** 引擎缓存上限(超过后淘汰最旧的非运行中引擎) */ const MAX_ENGINES = 30; export class AgentEngineManager extends EventEmitter { private engines = new Map(); /** 运行中的会话(stateChange INIT 添加 / TERMINATED 移除),用于 LRU 淘汰保护 */ private running = new Set(); /** 主 adapter(供 MemoryConsolidator 等共享组件使用) */ private primaryAdapter: IMetonaProviderAdapter; private fallbackAdapter: IMetonaProviderAdapter | null = null; private baseConfig: Partial; private workspacePath = ''; constructor( private opts: { /** adapter 工厂(每次调用返回新实例;闭包内读取最新配置) */ buildAdapter: () => IMetonaProviderAdapter; baseConfig: Partial; toolRegistry?: ToolRegistry; preToolHooks?: PreToolHook[]; postToolHooks?: PostToolHook[]; }, ) { super(); this.primaryAdapter = opts.buildAdapter(); this.baseConfig = { ...opts.baseConfig }; } /** 主 adapter(供 consolidator / orchestrator 等共享使用) */ getAdapter(): IMetonaProviderAdapter { return this.primaryAdapter; } /** 创建独立 adapter 实例(每引擎/SubAgent 独享,避免 abort 信号互踩) */ createAdapter(): IMetonaProviderAdapter { return this.opts.buildAdapter(); } getFallbackAdapter(): IMetonaProviderAdapter | null { return this.fallbackAdapter; } /** 设置故障转移 Provider(同步到所有引擎) */ setFallbackAdapter(adapter: IMetonaProviderAdapter | null): void { this.fallbackAdapter = adapter; for (const engine of this.engines.values()) { engine.setFallbackAdapter(adapter); } } getWorkspacePath(): string { return this.workspacePath; } setWorkspacePath(path: string): void { this.workspacePath = path; for (const engine of this.engines.values()) { engine.setWorkspacePath(path); } } /** 同步工具列表到所有引擎(工具开关变更 / MCP 注册完成时) */ setToolsAll(tools: MetonaToolDef[]): void { for (const engine of this.engines.values()) { engine.setTools(tools); } } /** 热更新所有引擎配置(设置变更时) */ updateConfigAll(partial: Partial): void { this.baseConfig = { ...this.baseConfig, ...partial }; for (const engine of this.engines.values()) { engine.updateConfig(partial); } } /** * 重建所有 adapter(LLM 配置变更时由 reloadAdapter 调用) * 工厂闭包读取最新配置,重建 primary + 各引擎独立实例。 */ refreshAdapters(): void { this.primaryAdapter = this.opts.buildAdapter(); for (const engine of this.engines.values()) { engine.setAdapter(this.createAdapter()); engine.setFallbackAdapter(this.fallbackAdapter); } } /** 获取(或创建)会话引擎 */ getEngine(sessionId: string): AgentLoopEngine { let engine = this.engines.get(sessionId); if (!engine) { engine = this.createEngine(sessionId); this.engines.set(sessionId, engine); this.evictIdleEngines(); } return engine; } /** * v0.7.4 P1-3: 会话当前是否有 run 进行中。 * 供 agent:sendMessage 做同会话并发防重入(第二个 invoke 在第一个 run 未结束时 * 直接拒绝,而非让同一引擎并行 runStream 或排队 30s 后强制 abort)。 */ isRunning(sessionId: string): boolean { return this.running.has(sessionId); } /** * v0.7.4 P3-5: 会话删除/归档时显式销毁引擎(内存收口)。 * 旧实现引擎只靠 LRU 上限 30 淘汰,会话删除后引擎与 adapter 实例仍驻留内存。 * 调用方在 sessions:delete 时联动调用;正在运行中的会话由调用方先 abort。 */ disposeEngine(sessionId: string): void { const engine = this.engines.get(sessionId); if (!engine) return; engine.destroy(); this.engines.delete(sessionId); this.running.delete(sessionId); log.debug(`[EngineManager] disposed engine for session ${sessionId}`); } /** 中断指定会话引擎 */ abort(sessionId: string): void { this.engines.get(sessionId)?.abort(); } /** 等待指定会话当前 run 结束 */ async waitForAbort(sessionId: string, timeoutMs = 5_000): Promise { return this.engines.get(sessionId)?.waitForAbort(timeoutMs) ?? true; } /** 当前引擎数量(测试/诊断用) */ get size(): number { return this.engines.size; } // ===== 私有方法 ===== private createEngine(sessionId: string): AgentLoopEngine { const engine = new AgentLoopEngine( { ...this.baseConfig, contextWindow: this.primaryAdapter.getContextWindow(), }, this.createAdapter(), this.opts.toolRegistry, this.opts.preToolHooks ?? [], this.opts.postToolHooks ?? [], ); engine.setFallbackAdapter(this.fallbackAdapter); if (this.workspacePath) engine.setWorkspacePath(this.workspacePath); // v0.5.2 关键修复: 新建引擎从 registry 拉取当前启用工具。 // setToolsAll 只同步"已存在"的引擎 — 引擎是懒创建的(首次 sendMessage 时 getEngine), // 启动期的 setToolsAll 调用时 engines Map 为空,全是 no-op。 // 此前缺此调用 → 新引擎 this.tools=[] → LLM 请求不带 tools → // 模型无法发起 tool_call(症状:模型口头说要调工具,实际不调,凭历史记忆瞎编)。 // v0.4.0 P2-10 引入每会话引擎时遗留的回归,v0.5.2 修复。 if (this.opts.toolRegistry) { engine.setTools(this.opts.toolRegistry.listTools()); } this.forwardEngineEvents(engine, sessionId); log.debug( `[EngineManager] engine created for session ${sessionId} (total: ${this.engines.size + 1})`, ); return engine; } /** 将引擎事件转发到 manager(统一附加 sessionId,供常驻监听器消费) */ private forwardEngineEvents(engine: AgentLoopEngine, sessionId: string): void { engine.on('streamEvent', (event) => { this.emit('streamEvent', { ...event, sessionId: event.sessionId || sessionId }); }); engine.on('stateChange', (data) => { const payload = { ...data, sessionId: data.sessionId || sessionId }; // 维护运行中集合(LRU 淘汰保护) if (payload.state === 'INIT' || payload.current === 'INIT') this.running.add(sessionId); if (payload.state === 'TERMINATED' || payload.current === 'TERMINATED') this.running.delete(sessionId); this.emit('stateChange', payload); }); engine.on('complete', (data) => { this.running.delete(sessionId); this.emit('complete', { ...data, sessionId: data.sessionId || sessionId }); }); engine.on('compressed', (data) => { this.emit('compressed', { ...data, sessionId }); }); engine.on('deadLoop', (data) => { this.emit('deadLoop', { ...data, sessionId: data.sessionId || sessionId }); }); engine.on('providerSwitched', (data) => { this.emit('providerSwitched', { ...data, sessionId: data.sessionId || sessionId }); }); engine.on('aborted', () => { this.emit('aborted', { sessionId }); }); } /** LRU 淘汰:超过上限时删除最旧的非运行中引擎 */ private evictIdleEngines(): void { if (this.engines.size <= MAX_ENGINES) return; let toEvict = this.engines.size - MAX_ENGINES; for (const [sessionId, engine] of this.engines) { if (toEvict <= 0) break; if (this.running.has(sessionId)) continue; engine.destroy(); this.engines.delete(sessionId); log.debug(`[EngineManager] evicted idle engine for session ${sessionId}`); toEvict--; } } }