框架能力

事件系统

探索 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

typescript
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`);
    });
}

步骤二:在自定义插件中主动发射事件

当你在编写自研的大模型适配器或工作流插件时,也可以向微内核事件总线主动发射事件:

typescript
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 进行收敛管理。