插件解耦神器!DeepSeek Harness 五种事件分发模式

📅 2026/8/20 10:04:50
插件解耦神器!DeepSeek Harness 五种事件分发模式
通过前面几篇我们已经知道插件怎么摆到桌子上也知道装配顺序怎么推导。剩下一个关键问题插件之间怎么互相通信想在LLM 调用里加一层重试 / 记账 / 拦截代码应该长什么样为什么waterfall是整个架构的枢纽答案就是 Cordis 的五种事件分发模式emit / parallel / serial / bail / waterfall以及基于它们做出来的一整套类型化事件系统。1. 五种分发模式一张表看清vendor/cordis/src/events.ts:32exporttypeDispatchModeemit|parallel|serial|bail|waterfall模式是否await顺序返回值典型用法emit否并发触发不 await忽略“已经发生了什么” —— 通知类事件parallel是并发 allSettled聚合错误多个副作用并行多路持久化serial是顺序await首个 bail 值需要 await 的顺序处理 短路bail否同步顺序首个 bail 值“谁能接就谁接”、拦截替换waterfall视监听器而定洋葱式包裹next最外层监听器返回around-middleware拦截改写转发Cordis 的判定是否算 bail很简单vendor/cordis/src/events.ts:13exportfunctionisBailed(value:any){returnvalue!nullvalue!falsevalue!undefined}null / false / undefined都表示我不接手其它值都表示我截胡了别再往下传。1.1emit一次广播不等结果vendor/cordis/src/events.ts:194emit(...args:any[]){this.dispatch(emit,args).map(cbcb(...args))}语义所有 listener 同时被调用一次不看它们返回什么、也不 await Promise。这是通知类事件最合适的模式。dsh 里的例子session/message-appended、tool/result—— 只是把发生了什么广播出去让日志/UI/统计各自消费。1.2parallel并发 await 聚合错误vendor/cordis/src/events.ts:183asyncparallel(...args:any[]){constresultsawaitPromise.allSettled(this.dispatch(emit,args).map(asynccbcb(...args)))consterrorsresults.filter((r):risPromiseRejectedResultr.statusrejected)if(errors.length)thrownewAggregateError(errors.map(errorerror.reason))}语义所有 listener 同时被调用allSettled等它们全部结束有失败就聚合成AggregateError一起抛。适合若干条持久化通道并发写入我要等所有都写完其中一条挂了也不影响别的开始。1.3serial顺序 await遇到 bail 就返回vendor/cordis/src/events.ts:204asyncserial(...args:any[]){for(constcbofthis.dispatch(serial,args)){constresultawaitcb(...args)if(isBailed(result))returnresult}}语义一个一个await哪个 listener 返回非 nullish、非 false就把它作为整体结果吐回去、别的都不跑了。用来做分层处理链session/attach-context之类——每一层看看要不要挂点上下文处理不了就返回undefined让下一层来。1.4bail同步版本的 serialvendor/cordis/src/events.ts:217bail(...args:any[]){for(constcbofthis.dispatch(bail,args)){constresultcb(...args)if(isBailed(result))returnresult}}语义跟serial一样但不 await。适合同步的谁能处理谁接手比如注册钩子时的策略选择“谁给我一个替代方案”。Cordis 内部的internal/listener就是bail// vendor/cordis/src/events.ts:296constresultthis.bail(this.ctx,internal/listener,name,listener,options)if(result)returnresult如果某个 listener 说这个事件我要特殊处理就直接用它的返回值代替常规注册流程。1.5waterfall洋葱模型重点这是 dsh 里出场率最高、分享要花时间讲的模式。见下节。2.waterfall内部洋葱模型是怎么实现的只看核心代码vendor/cordis/src/events.ts:234-243waterfall(...args:any[]){constcbsthis.dispatch(waterfall,args)// 拿到所有 listenerconstinnerargs.pop()// 参数列表最后一位 原始行为constnext(){constcbcbs.shift()??inner// 没 listener 了 → 跑原始行为returncb(...args)// 传给下一层的还是同一份 args}args.push(next)// 把 next 塞回参数末尾returnnext()// 从最外层 listener 起跑}signature 是(...args, next)每个 listener 收到的最后一个参数就是next调next()→ 把执行权交给下一层 listener或最内层的原始行为不调next()→ 短路整个链条到你为止用你的返回值作为最终结果这是AOP 环绕通知、“中间件”、onion model的通用模式listener A (outermost) ├─ 前置改 args、加超时、记账开始 ├─ const result next() ← 委派进 B │ listener B │ ├─ 前置 │ ├─ const result next() ← 委派进 C │ │ listener C │ │ ├─ 前置 │ │ ├─ const result next() ← 已经没 listener 了 → 跑 inner │ │ │ inner (原始行为) │ │ ├─ 后置包装 result │ │ └─ return │ ├─ 后置 │ └─ return ├─ 后置wrap for retry └─ return ← 最外层返回值 waterfall 的结果dsh 的硬约束packages/CLAUDE.mdWaterfall listeners MUST callnext()to delegate; returning without it short-circuits the chain.“忘记调next()是新人最常踩的坑”。表现是某个 listener 加进去之后后面所有的行为都消失了比如整个 LLM 调用不动了。原因就是短路了。3. LLM 调用llm/streamwaterfall来看真正跑生产的例子packages/llm/llm/src/index.ts:917privatestreamWithRegistration(options:GenerateOptions,prepared?:{registration:AdapterRegistration;config:LlmCallConfig},):AsyncIterableStreamChunk{returnthis.ctx.waterfall(this,// thisArg LlmRuntime 实例llm/stream,// 事件名options,// 参数调用选项()this.adapterStream(options,prepared),// ★ inner真实发 HTTP 调 DeepSeek)}三层含义this LlmRuntime 实例—— listener 里可以this.registerAdapter(...)、this.adapters一样使用。参数是options—— listener 收到(options, next)可以任意读改。inner是() this.adapterStream(...)—— 就是没人拦截时应该发生的行为实际上会挑选 provider 的 adapter然后调它的.stream()。任何第三方插件都可以往这里挂拦截// 一个虚构的 retry 插件exportconstnamellm-retryexportconstinject[llm]exportfunctionapply(ctx:Context){ctx.on(llm/stream,function(options,next){// ── 前置可以改 options、加 abort signal、记账ctx.logger.info([llm] provider%s model%s,options.provider,options.model)// ── 委派走原始 stream 或再下一层 listenerconststreamnext()// ── 后置给流套一层遇到错就重试returnwrapForRetry(stream,options,ctx)})}这条设计打通了几件事重试——llm-retry插件现网真实存在就是这么写的模型路由—— 一个 listener 根据options.model判断改写options.provider然后next()重放 / 录制—— snapshot 测试通过挂一个录制流的 listener 把 chunk 存下来计费 / 归因——llm/attribution相关插件在 listener 里把 tokens 加到账本一条硬约束来自 00 · 索引Model-visible ⟺ Logged——所有能到达模型请求的输入都必须能从 session log 完整重放。挂载llm/streamlistener 时如果偷偷塞 token 或改 prompt会破坏这条不变量。因此**llm/stream只允许改传输/调度相关的东西不允许无中生有加内容**。4. 工具执行三段 waterfall 流水线packages/core/tools/src/index.ts把一次工具调用拆成了三段独立的 waterfallpackages/core/tools/src/index.ts:1475-1746// 第一段pre-execute —— 审批 / 鉴权 / 参数改写constgateawaitthis.ctx.waterfall(carrier,tools/pre-execute,exec,/* inner */...)// 第二段execute —— 真正跑工具constresultawaitthis.ctx.waterfall(carrier,tools/execute,mutableExec,/* inner */...)// 第三段post-execute —— 结果加工 / Loop hygiene / 二次校验constdecisionawaitthis.ctx.waterfall(scopeTarget(this,exec.agent),tools/post-execute,exec,result,/* inner */...)三个事件的类型签名packages/core/tools/src/index.ts:152-175interfaceEvents{/** mode waterfall */tools/pre-execute(this:ScopedToolRuntime,exec:ToolExecution,next:()PromisePreToolDecision):PromisePreToolDecision/** mode waterfall */tools/execute(this:ScopedToolRuntime,exec:ToolDispatchExecution,next:()PromiseToolExecutionResult):PromiseToolExecutionResult/** mode waterfall */tools/post-execute(this:ScopedToolRuntime,exec:ToolExecution,result:ReadonlyToolExecutionResult,next:()PromisePostToolDecision):PromisePostToolDecision}为什么要拆成三段因为语义不同pre-execute结果是PreToolDecision允许/拒绝/替换/等审批execute结果是ToolExecutionResult工具真跑完了post-execute结果是PostToolDecision结果放行/替换/终止 loop不同的插件挂不同的段插件挂载点目的interactiontools/pre-execute审批把是否允许执行抛给用户guard(tool-call-timeout)tools/execute超时包裹执行加 abortguard(loop-hygiene)tools/post-execute循环治理连续几次同工具调用就打断session三段都挂把每次 tool call 记入 log互不感知interaction不需要知道guard存在session不需要知道interaction存在。它们只是各自挂在自己关心的 waterfall 上都调next()把执行权交回主流水线。5. 类型化事件签名走声明合并除了运行时行为Cordis 事件的类型系统也是一大亮点。回到packages/llm/llm/src/index.ts:46declaremoduledeepseek-ai/cordis{interfaceContext{llm:LlmRuntime}interfaceEvents{/** * Waterfall around every streaming model call (retry, replay, routing). * mode waterfall */llm/stream(this:LlmRuntime,options:GenerateOptions,next:()AsyncIterableStreamChunk,):AsyncIterableStreamChunk}}Events是一个可扩展的接口vendor/cordis/src/events.ts:329exportinterfaceEvents{internal/plugin(fiber:Fiber):voidinternal/status(fiber:Fiber,oldValue:FiberState):void// …只有框架自带的几个 internal/* 事件}任何包都能通过declare module往这张表里补事件。每个事件名只在一处声明参数类型this类型ScopedToolRuntime/LlmRuntime/ …返回类型一次写完所有地方ctx.on、ctx.emit、ctx.waterfall、ctx.serial…都能拿到静态签名。5.1ctx.on也在 Context 接口里vendor/cordis/src/events.ts:97declaremodule./context.ts{exportinterfaceContext{onKextendskeyofEvents(name:K,listener:Events[K],options?:boolean|EventOptions,):()boolean// …parallel / emit / serial / bail / waterfall 都同样声明}}关键是listener: Events[K]—— 你写错事件名或者写错参数类型编译期就报错// ✅ OKctx.on(llm/stream,(options,next)next())// ❌ TS Error: 参数类型不匹配ctx.on(llm/stream,(options:string,next)next())// ❌ TS Error: 事件名不存在ctx.on(llm/stream-nonexistent,(){})这就是 “typed event” 的全部含义事件不是emit(string, any)是有静态签名的调用协议。5.2modeJSDoc 与 verify-event-modes回头看每个事件定义上面的mode注解/** * mode waterfall */llm/stream(...):AsyncIterableStreamChunk这不是装饰——scripts/verify-event-modes.ts会扫所有interface Events声明确保每个事件都有mode事件签名和 mode 匹配waterfall 必须带nextserial 返回可 bail 值等等事件调用点用的分发方法与 mode 一致不能ctx.emit(llm/stream, ...)这是让通信协议从约定变成编译期强制的关键设计。6. 事件监听器与 fiber 生命周期绑定这里是 05 的引子为什么插件卸载时 listener 会自动摘除看vendor/cordis/src/events.ts:254register(label:string,hooks:Hook[],callback:any,options:EventOptions):()void{constmethodoptions.prepend?unshift:pushreturnthis.ctx.fiber.effect((){// ★ 走 fiber.effecthooks[method]({ctx:this.ctx,callback,...options})// setup: 把 listener 塞进 hooks 数组return()this.unregister(hooks,callback)// teardown: 从 hooks 数组里摘走},label)}注册这一步登记为一个fiber.effect卸载这一步从数组里删掉 listener插件 fiber 从 ACTIVE 转 DISPOSED比如热更、比如依赖挂了框架会自动跑 teardown你不用写一个字的off()这就是注册即副作用、副作用可逆——ctx.on只是这个规律的一个具体实例。下一篇会集中讲这条主线。7. 内建的 hookdsh 用得上的几个internal/*事件Cordis 自己也走这套通信机制。看vendor/cordis/src/events.ts:329exportinterfaceEvents{internal/plugin(fiber:Fiber):void// 插件加载/卸载internal/status(fiber:Fiber,oldValue:FiberState):void// fiber 状态变了internal/config(this:Fiber,config:any,next:()any):any// waterfall解析 configinternal/service(this:Context,name:string,value:any):void// 服务绑定internal/update(this:Fiber,config:any,noSave:boolean,next):void|Promisevoid// waterfall更新 configinternal/get(ctx:Context,name:string,error:Error,next:()any):any// waterfall读服务internal/set(ctx:Context,name:string,value:any,error:Error,next):boolean// waterfall写服务internal/listener(this:Context,name:string,listener:any,prepend:boolean):void// bail注册 listenerinternal/dispatch(mode:DispatchMode,name:string,args:any[],thisArg:any):void// 有事件分发发生了}dsh 用得上的internal/status—— 监控某个 fiber 从 PENDING 到 ACTIVE 的变化用于诊断 / 状态灯internal/dispatch—— 全局旁路看到所有非 internal 事件被派发的现场做 tracing / metrics 很方便internal/plugin—— 观察所有插件的挂载/卸载当你想做跨插件观察而不是参与时这些内建事件就是入口。8. 五种模式怎么挑一句话决策树“我要广播一个事件我该用哪个”一个简单的决策树你需要 listener 的返回值 / 想让某个 listener 短路整条链吗 ├─ 是 → 你需要 await Promise 吗 │ ├─ 是 → 用 serial要 await 短路 │ ├─ 否 → 用 bail同步 短路 │ └─ 你想让 listener 环绕原始行为前置/委派/后置→ 用 waterfall └─ 否 → 你需要 await Promise 吗 ├─ 是 → 用 parallel并发 await 完成 └─ 否 → 用 emitfire-and-forgetwaterfall 的语义特别—— 它不是顺序/并发 bail的哪个变体而是环绕能同时做前置、后置、否决。挑不出模式时优先考虑waterfall。9. 代码位置速查主题文件关键位置DispatchMode枚举vendor/cordis/src/events.tsL32isBailed判定vendor/cordis/src/events.tsL13Context 接口扩展on / emit / serial / …vendor/cordis/src/events.tsL34-109EventsService.dispatchvendor/cordis/src/events.tsL165-175parallel实现vendor/cordis/src/events.tsL183-187emit实现vendor/cordis/src/events.tsL194-196serial实现vendor/cordis/src/events.tsL204-209bail实现vendor/cordis/src/events.tsL217-222waterfall实现vendor/cordis/src/events.tsL234-243register走 fiber.effectvendor/cordis/src/events.tsL254-260ctx.on逻辑vendor/cordis/src/events.tsL288-302内建Events定义vendor/cordis/src/events.tsL329-352llm/stream声明packages/llm/llm/src/index.tsL46-67llm/streamwaterfall 调用点packages/llm/llm/src/index.tsL917-927工具三段 waterfall 声明packages/core/tools/src/index.tsL150-175工具 pre/execute/post 调用点packages/core/tools/src/index.tsL1475 / L1573 / L1743事件模式静态检查scripts/verify-event-modes.ts全文硬约束waterfall 必须调 next()packages/CLAUDE.md“Waterfall listeners MUST callnext()” 段