事件系统
探索 Cordis 异步事件总线、中间件拦截器(Bail/Waterfall)与 Harness 核心事件分发拓扑。
在智能体的自主推理与工具执行过程中,系统内部会发生高频的状态流转——从会话开始、Prompt 组装、大模型流式推理,到工具拦截与 Trajectory 轨迹落盘。
传统的全局单例 EventEmitter 在插件频繁热重载时极易引发事件重复触发与严重的内存泄漏。DeepSeek Harness 基于 Cordis 微内核构建了具备作用域追踪、短路守卫拦截(Bail) 与 顺序流水线变换(Waterfall) 的企业级事件总线。
本文将以一个 “智能体高危命令拦截与提示词敏感词过滤探针(Command Security Interceptor)” 为例,系统剖析 Harness 事件体系的三大调度模式。
核心架构:Cordis 事件体系的三大调度模式
Cordis 事件总线根据业务语义,提供了三种截然不同的分发机制:
┌─────────────────────────────────────────────────────────────┐
│ Cordis 微内核三大事件调度模式 │
├─────────────────────┬───────────────────────────────────────┤
│ 1. emit (广播通知) │ 并发通知所有监听器,不关心返回值 │
│ │ 典型场景:日志记录、指标上报、审计追溯│
├─────────────────────┼───────────────────────────────────────┤
│ 2. bail (短路拦截) │ 按序触发监听器,遇到首个非 undefined │
│ │ 返回值立即中断后续执行(安全熔断守卫)│
├─────────────────────┼───────────────────────────────────────┤
│ 3. waterfall (水流) │ 类似流式管道,前一个监听器的输出作为 │
│ │ 后一个监听器的输入(Prompt 过滤脱敏) │
└─────────────────────┴───────────────────────────────────────┘
步骤一:实战编写安全拦截器插件
创建 custom-plugins/security-interceptor/src/index.ts:
import type { Context } from '@deepseek-ai/cordis';
export const name = 'security-interceptor-plugin';
export function apply(ctx: Context) {
// 1. 【Bail 短路拦截】:在工具真正执行前进行安全审查
ctx.on('tool/before-execute', (call) => {
// 如果检测到高危命令,直接返回错误信息,阻断后续工具执行
if (call.name === 'bash' && call.args?.command?.includes('rm -rf /')) {
console.warn(`[Security] 拦截到高危命令调用: ${call.args.command}`);
return {
blocked: true,
reason: '安全策略拒绝执行破坏性全局删除指令'
};
}
// 返回 undefined 表示放行,微内核继续执行下一个监听器
});
// 2. 【Waterfall 管道变换】:在发送给模型前脱敏 Prompt 文本
ctx.on('prompt/transform', (inputPrompt: string) => {
// 将明文手机号替换为脱敏掩码
const sanitized = inputPrompt.replace(/(\d{3})\d{4}(\d{4})/g, '$1****$2');
return sanitized;
});
// 3. 【Emit 广播通知】:无阻塞异步记录审计日志
ctx.on('trajectory/step', (step) => {
console.log(`[AuditLog] 实时落盘会话步骤 #${step.index},耗时: ${step.durationMs}ms`);
});
}
步骤二:在自定义插件中主动发射事件
当你在编写自研的大模型适配器或工作流插件时,也可以向微内核事件总线主动发射事件:
import type { Context } from '@deepseek-ai/cordis';
export function runWorkflow(ctx: Context, taskPayload: any) {
// 1. 发射短路检查事件:如果有任何安全插件阻断,立即中止工作流
const guardResult = ctx.bail('workflow/before-run', taskPayload);
if (guardResult?.blocked) {
console.error(`[Workflow] 任务被安全插件拦截: ${guardResult.reason}`);
return;
}
// 2. 发射流水线变换事件:依次通过所有过滤插件处理 payload
const finalPayload = ctx.waterfall('workflow/transform-payload', taskPayload);
// 3. 广播通知执行成功
ctx.emit('workflow/completed', { id: taskPayload.id, timestamp: Date.now() });
}
Harness 内置核心生命周期事件速查
┌─────────────────────────────────────────────────────────────┐
│ Harness 核心生命周期事件流 │
├─────────────────────────────────────────────────────────────┤
│ • session/before-turn : 每一轮智能体推理前触发 │
│ • prompt/transform : 发送给 LLM 前的 Prompt 过滤流 │
│ • llm/stream-chunk : 模型输出 Token 流式分片到达 │
│ • tool/before-execute : 工具执行前的安全拦截检测 │
│ • tool/after-execute : 工具完成并生成观测数据后触发 │
│ • session/after-turn : 当前轮次推理完成并落盘轨迹 │
└─────────────────────────────────────────────────────────────┘
常见问题解答 (FAQ)
Q1: 为什么 Cordis 事件监听器不需要手动调用 removeListener?
解答:因为 ctx.on 是在当前插件专属的 Context 作用域内注册的。当插件被卸载或热重载时,微内核会自动解绑该上下文下所有已挂载的监听器,彻底根治 Node.js 经典的 MaxListenersExceededWarning 内存泄漏问题。
Q2: ctx.bail 触发时,监听器的执行先后顺序是怎样的?
解答:按照插件加载激活的时间先后顺序执行。一旦某个监听器返回了非 undefined 值,微内核将立即“短路(Short-Circuit)”,后续排队的监听器将不再被调用。
Q3: ctx.waterfall 与普通数组 reduce 有什么区别?
解答:ctx.waterfall 完美支持异步 Promise 链式等待。上一个异步插件处理完数据并返回 Promise 后,下一个监听器会自动等待该 Promise 解析后的新数据作为入参,非常适合异步多级提示词流水线。
Q4: 可以在事件监听器内部再次调用 ctx.emit 触发其他事件吗?
解答:完全可以。但在设计事件联动逻辑时需注意避免“事件乒乓(Event Ping-Pong)”形成无限递归调用链。建议在复杂的业务状态流转中使用状态机或 Service 进行收敛管理。