Skip to content

05 事件系统:DSH 的扩展点

事件是 DSH 的第一扩展点:不改动循环本身的一行代码,只靠监听事件就能记录、改写、拦截甚至否决模型轮次中的任意环节——DSH 的微内核主张「没有任何一行修改循环本身」,正是靠这张事件网实现的。

本章目标

  • 理解事件三域并会选域;
  • 吃透 waterfall 的授权语义(reject / 改写 / enter(messages),必须调 next());
  • agent/pre-step 记录、改写模型即将看到的消息;
  • tools/pre-execute / tools/post-execute 拦截工具调用;
  • agent/turn-stopping 做轮次终点检查点、用 session/event 观察持久事件流;
  • 记住监听器纪律。

事件三域:先决定挂在哪里

选对事件域是大多数改动的第一个决定:

代表事件性质什么时候用
session 事件turn/startstep/startuser/messageassistant/chunkassistant/messagetool/calltool/resultturn/end追加进只增日志的持久事实,统一经 session/event 广播需要重载后仍存在、可回放、可重建的事实
agent 事件agent/*agent/statusagent/pre-stepagent/requestagent/turn-stoppingagent/inbox/*实时协调,payload 携带活跃的 live Agent观察或拦截进行中的工作
capability 事件fs/*tools/*telemetry/*挂到 seam 上的策略 / 适配器,无需导入循环给某项能力附加策略

判定规则:要持久 → session 事件;要实时干预agent/*要拦能力 → capability 事件。另有一个贯穿全程的不变量:Model-visible ⟺ logged(模型可见即已记录)——一切抵达模型请求的内容都必须能从会话日志重建。

五种分发模式回顾:waterfall 是授权性拦截点

模式语义监听器返回值
emit同步广播,不等待、不收集忽略
parallel并发运行并一同等待忽略
serial按序 await,遇到首个 bail 值(非 null/false/undefined)提前终止首个 bail 值胜出
bailserial 的同步版本同 serial
waterfall环绕中间件:监听器收 (...args, next)next() 委托给下游最外层监听器的返回值

waterfall 是授权性拦截点:调用 next() 会执行下一个监听器,最终落到内置默认行为;不调用 next() 就是否决(veto),短路整条链。监听器可以:① 观察并委托(记录后 return next());② 改写后委托(改共享 payload 再 return next(...));③ 直接裁决(不调 next(),返回自己的决策)。agent/pre-stepagent/requestllm/streamtools/pre-executetools/post-execute 都是 waterfall。纪律只有一条:只观察/标注的监听器必须调用 next()——忘了调用,下游默认行为会被吞掉。

实战一:agent/pre-step —— 决定模型看到什么

agent/pre-step 在模型请求被构造之前触发,是「模型将看到什么」的唯一权威入口。payload 携带独占的已领取消息批次,监听器要么拒绝整个提议的步骤,要么决定以哪组消息进入:

ts
'agent/pre-step'(this: Scoped<Agent>, payload: {
  agent: Agent
  messages: UserMessage[]   // 独占的已领取消息批次
  turn: number
  step: number
  signal: AbortSignal
}, next: () => Promise<PreStepDecision>): Promise<PreStepDecision>
// PreStepDecision = { kind: 'reject' } | { kind: 'enter'; messages: UserMessage[] }

先写一个只记录、不改写的监听器——它必须 return next(),否则模型什么都看不到:

ts
import type { Context } from '@deepseek-ai/cordis'

export const name = 'pre-step-audit'

export function apply(ctx: Context) {
  ctx.on('agent/pre-step', async (payload, next) => {
    const { messages, turn, step } = payload
    // 记录:这一步骤模型即将看到几条消息
    ctx.logger('pre-step').info(`turn ${turn} step ${step}: ${messages.length} message(s)`)
    // 观察型监听器必须委托给下游,否则会吞掉整个步骤
    return next()
  })
}

改写示例:在消息尾部注入一条提醒,把改写后的消息通过 enter 交还:

ts
import { createUserMessage } from '@deepseek-ai/dsh-llm'

ctx.on('agent/pre-step', async ({ messages }, next) => {
  const reminder = createUserMessage({
    content: [{ type: 'text', text: '(注入)回答前先检查 todo 状态。' }],
    source: { kind: 'inject' }, // inject 来源:不唤醒 driver,只进下一次请求
  })
  return next({ kind: 'enter', messages: [...messages, reminder] })
})

也可以直接拒绝:返回 { kind: 'reject' } 会关闭一个不含步骤的持久轮次(日志仍会记录这次尝试)。agent/* 事件全部是 scope-filtered(this: Scoped<Agent>),只收到自己作用域内的调用。以返回的决策为准——包装 next() 的监听器保留下游消息,除非有意替换。

实战二:tools/pre-execute 与 tools/post-execute —— 拦截工具调用

工具执行流水线也是 waterfall 的天下,顺序是 tools/pre-execute → 单调 guard → tools/executetools/post-execute → 仅观测的 tools/result

tools/pre-execute 是调用前的门禁:返回 allow | deny | ask。它有意不允许改写 exec.arguments——参数已经记录并呈现,改写会造成日志与实际运行失同步:

ts
import type { Context } from '@deepseek-ai/cordis'
import type { PreToolDecision, ToolExecution } from '@deepseek-ai/dsh-tools'

declare function isAllowed(exec: ToolExecution): Promise<boolean>

export function apply(ctx: Context) {
  ctx.on('tools/pre-execute', async (exec, next): Promise<PreToolDecision> => {
    if (exec.name === 'bash' && !(await isAllowed(exec))) {
      // 不调 next():短路,直接否决这次调用
      return { kind: 'deny', reason: 'Denied by policy.' }
    }
    return next() // 放行
  })
}

tools/post-execute 在结果产生后运行:接受、替换某一投影(content 或 value 二选一)、附加 context,或用 block 把纠正性反馈变成错误结果。例如截断超长输出:

ts
import type { PostToolDecision } from '@deepseek-ai/dsh-tools'

ctx.on('tools/post-execute', async (exec, result, next): Promise<PostToolDecision> => {
  if (!result.isError && exec.name === 'read' && typeof result.value === 'string'
      && result.value.length > 4000) {
    // 替换 value 投影:保留前缀,明确标注截断
    return { kind: 'accept', value: result.value.slice(0, 4000) + '\n…(truncated)' }
  }
  return next()
})

实战三:agent/turn-stopping —— 轮次终点检查点

agent/turn-stoppingserial 事件,没有 next()。它在轮次自然停止(无工具、无 steering 后续)时、最后一次 steering 排空之前运行,是「这轮该不该结束」的终端检查点——监听器可以通过 agent.steer() 反对关闭:

ts
ctx.on('agent/turn-stopping', ({ agent }) => {
  if (goalStillOpen(agent)) {
    // 反对关闭:排队下一次迭代,轮次继续推进
    agent.steer(createUserMessage({
      content: [{ type: 'text', text: '继续:目标尚未完成。' }],
      source: { kind: 'inject' },
    }))
  }
})

serial 语义:按注册顺序依次 await,第一个返回 bail 值就提前终止——不过 turn-stopping 约定返回 void,它的「动作」是 steer 的副作用而非返回值。这类终点检查点常用于 hook 系统、目标驱动循环等「每轮都要看一眼」的策略。

观察持久事件流:session/event

sessions 把每个持久会话事件追加进只增日志,并同步广播为 session/event(emit 模式,忽略监听器返回值)。它是监听方最多的事件,UI、遥测、回放、transcript 都从这里渲染:

ts
ctx.on('session/event', (_session, event) => {
  if (event.type === 'assistant/chunk' && event.data.chunk.type === 'text-delta') {
    render(event.data.chunk.text) // 逐块渲染助手输出
  }
})

高频事件速查

事件模式用途
session/eventemit持久会话事实流:UI、回放、遥测、标题生成
agent/statusemitidle / running 实时状态,驱动状态显示与自动化区间
agent/pre-stepwaterfall决定模型看到什么:记录、改写、拒绝
agent/requestwaterfall替换冻结的模型调用配置(不能改消息——模型可见内容必须走已记录通道)
tools/pre-executewaterfall工具调用前门禁:allow / deny / ask
tools/post-executewaterfall工具结果后处理:截断、替换、附加 context

(还有 llm/stream(waterfall,流式输出改造)、agent/turn-stopping(serial)等,完整矩阵见事件生产方与消费方文档。)

监听器纪律

  1. 回调异常隔离:监听器抛异常不得 reject 所在 promise,也不得饿死后续监听器——分发器会 try/catch 并记日志。业务失败用返回值/决策表达,不要抛异常。
  2. 注册返回 disposerctx.on(...) 返回移除函数;监听器归当前 fiber 所有,插件卸载(含热重载)时自动移除。
  3. waterfall 短路语义:不调 next() 等于有意否决;只观察的监听器必须 next()。协作式监听器常改写共享 payload(如 messages)后再委托,让下游看到改动后的值。

小结

事件就是扩展点:持久事实session/event实时协调agent/*(携带 live Agent),策略 / 适配器挂 capability seam。waterfall 是授权性拦截点,serial 是终端检查点,emit 是广播。写监听器时记住:观察者必须 next()、注册自带 disposer、业务失败用返回值表达。

下一步

下一章可以继续两条线:要么看 llm/stream 与 LLM 适配器,在流式输出层改造模型行为;要么回到工具注册,把 tools/pre-execute 的门禁与审批(ctx.approval)组合成完整的权限系统。