Skip to content

EventBus 与 HookRegistry — 事件与中间件

Aalis 提供两种互补的扩展机制:事件(单向通知)和中间件钩子(可拦截的管道)。

EventBus — 事件总线

源码: packages/core/src/primitives/events.ts

类型安全的全局发布/订阅事件总线,用于松耦合的异步通知。事件只通知、不干预流程。

API

typescript
// 监听事件(返回 dispose 函数)
const off = ctx.on('inbound:message', async (msg) => { ... });

// sticky 事件:注册晚于发出也能在下一个微任务收到补发('app:ready' / 'app:started')
ctx.on('app:ready', () => { ... });

// 发出事件(按注册顺序依次 await 每个 handler,永不拒绝)
await ctx.emit('outbound:message', outMsg);

实现特性

  • 异步串行:emit() 会 await 每个 handler 完成后再执行下一个
  • 按注册顺序调用
  • on() 返回 dispose 函数,可随时移除监听
  • 每次 on() 是一条独立登记:同一函数登记两次触发两次,各自的 dispose 只移除自己那条;两个 Context 共用同一函数也互不影响
  • Context 销毁时自动移除该 Context 注册的所有监听

core 内置事件

core 自持的只有下面十一个基础设施事件(源码 packages/core/src/types/events.ts)。消息、工具、会话、 调度等业务事件不在 core,由各 -api 包经 declaration merging 注入——谁注入了哪一族见 扩展点索引 §2,键与 payload 的权威定义在该包的 declare module 声明里。

十一个事件按「发射方等不等监听器」分两节,节由前缀判定:

屏障app:*):emit 是 App 生命周期方法里的一步,监听器全部返回后才推进下一步。

事件参数时机
app:startingstart() 的第一步
app:ready启动第一相位(sticky)
app:started启动第二相位:全部 app:ready 监听器完成之后(sticky);CLI / TUI 在此接管终端
app:restartingrestart() 先发本事件,监听器全部完成后才把控制交给宿主的 RestartStrategy
app:stoppingstop() 开头,插件拓扑逆序 dispose 之前

通知service:* / plugin:* / plugins:changed):发射方不等监听器。发射点要么是同步的注册 / 拆卸收尾, 要么在 PluginManager 的 recompute flight 或挂起段内——那里等监听器会与 plugins.idle() 互等死锁。 监听器因此不能假设「我返回了状态机才继续」;要看落定后的状态请 await plugins.idle()

事件参数时机
service:registeredname某服务多了一个提供者(ctx.provide
service:unregisteredname某服务少了一个提供者(退订闭包或 Context 拆卸)
service:preference-changedname该服务的偏好 provider 切换(preferService / unpreferService);whenService 借此重挂
plugin:loadedinstanceId插件实例已激活;同一轮 recompute 可能紧接着激活下一个插件
plugin:unloadedinstanceId插件实例已拆卸;激活失败的回滚与关机拆卸不发
plugins:changed一轮 recompute 收敛,插件状态集合可能已变;关机轮不发

app:readyapp:started 是两个相位,不是同一里程碑的两个名字:start() 串行 await,app:started 严格晚于 全部 app:ready 监听器完成。app:stopping 用于知会(打印告别语、切状态条),不是清理通道——清理副作用一律走 ctx.onDispose(fn),它覆盖 bounce / unload / 停机全部路径。总线上没有 dispose 事件。

扩展自定义事件

第三方插件通过 TypeScript declaration merging 即可为事件系统新增类型安全的自定义事件:

typescript
declare module '@aalis/core' {
  interface AalisEvents {
    'scheduler:tick': [jobId: string];
    'scheduler:error': [jobId: string, error: Error];
  }
}

// 之后可以类型安全地使用
ctx.on('scheduler:tick', async (jobId) => { ... });
ctx.emit('scheduler:tick', 'job-1');

AalisEvents类型封闭的(没有 [key: string] 兜底):对扩展开放、对拼写错误封闭—— 未声明的事件名在 ctx.on / ctx.emit 处直接编译报错,契约始终可枚举、依赖边在包图中可见。

事件名需要运行时动态生成(如按频道/任务 ID 派生)时,官方出路是在自己的命名空间内 合并一条模板字面量签名(TS 4.4+):

typescript
declare module '@aalis/core' {
  interface AalisEvents {
    // 动态事件名族:myplugin:channel: 前缀下的任意后缀都合法,payload 类型统一
    [k: `myplugin:channel:${string}`]: [payload: ChannelMessage];
  }
}

ctx.on(`myplugin:channel:${channelId}`, async (msg) => { ... }); // msg: ChannelMessage

同前缀下更具体的字面量 key 仍可逐条声明(TS 优先匹配字面量)。前缀必须用自己插件的 命名空间,避免与他人模板签名相互吞并。


HookRegistry — 中间件钩子管道

源码: packages/core/src/primitives/hooks.ts

中间件钩子是 Aalis 最强大的扩展机制。与事件不同,钩子是有序管道,插件可以修改管道中的数据、也可以完全中断流程。

核心概念

中间件采用 (data, next) 签名。调用 next() 将控制权传递给下一个中间件(或最终的 defaultAction)。不调用 next() 即中断整个管道——这是拦截消息的标准做法。

ctx.runHook(hookName, data, defaultAction)


中间件 A(先注册) ───── await fn(data, next)
  │ next()                     │ 不调用 next() → 中断
  ▼                             ▼
中间件 B(后注册)        管道终止,defaultAction 不执行
  │ next()

defaultAction() ← 所有中间件都 next() 后执行

API

typescript
// 注册中间件(同一钩子内按注册顺序执行)
const dispose = ctx.middleware('agent:reply:before', async (data, next) => {
  data.content = processContent(data.content);
  await next();
});

// 执行管道(由 Agent 或其他插件调用)
await ctx.runHook('agent:reply:before', { content: '...' }, async () => {
  // defaultAction: 所有中间件通过后才执行
});

// dispose() 可手动解除;插件卸载时本 ctx 注册的中间件自动清扫

内置钩子

core 的 HookContextMap 是空接口:内核不内置任何钩子键,全部由 -api 包注入。本仓库第一方包注入的键族如下, 键名与 data 类型的权威定义在各包的 declare module 声明里(按包查见扩展点索引 §3):

注入方钩子键用途
@aalis/api-agentagent:input:before / agent:llm:before / agent:llm:after / agent:tool:before / agent:tool:after / agent:reply:before / agent:turn:afteragent 一轮处理的各相位
@aalis/api-gatewayinbound:confirm / inbound:command / inbound:flow / inbound:trigger / inbound:dispatch / outbound:dispatch网关出入站的命名相位
@aalis/api-memorymemory:clear统一记忆清理编排,供 /clear 与各记忆插件协作

中间件特性

  • 注册顺序执行: 同一钩子内按注册顺序串行执行(无优先级数字;相位间次序由调度方显式表达)
  • 数据修改: data 通过引用传递,修改 data 对象即影响后续中间件和 defaultAction
  • 流程控制: 调用 next() 继续管道;不调用则中止后续中间件和 defaultAction
  • 上下文绑定: 每个中间件关联注册方 Context 的清理归属,插件卸载时自动清理(unregisterByOwner

典型用法

往提示词里加内容不要用本钩子——那是 agent:prompt 贡献点的活(ctx.contribute, 见 architecture.md 扩展机制):贡献点自带幂等、确定性排布与 错误隔离,中间件手搓 unshift 三者全无。本钩子留给改写 / 截停语义。

typescript
// 1. 拦截消息(不调用 next = 中断管道)
ctx.middleware('agent:input:before', async (data, next) => {
  if (shouldBlock(data.message)) return; // 不调用 next,整个管道终止
  await next();
});

// 2. 后处理回复内容
ctx.middleware('agent:reply:before', async (data, next) => {
  await next();
  data.content = transform(data.content);
});

// 3. 替换工具列表(如工具搜索层)
ctx.middleware('agent:llm:before', async (data, next) => {
  data.tools = await searchRelevantTools(data.messages);
  await next();
});

扩展自定义钩子

第三方插件可以定义自己的钩子,并让其他插件注入中间件:

typescript
// 声明类型(可选但推荐)
declare module '@aalis/core' {
  interface HookContextMap {
    'schedule:before': { jobId: string; cron: string };
  }
}

// 定义钩子的插件:在关键路径上调用 ctx.runHook
await ctx.runHook('schedule:before', { jobId, cron }, async () => {
  // defaultAction: 执行调度任务
  await executeJob(jobId);
});

// 拦截钩子的第三方插件
ctx.middleware('schedule:before', async (data, next) => {
  logger.info(`即将执行: ${data.jobId}`);
  data.cron = modifyCron(data.cron);
  await next();
});

自定义 hook 需要通过 declaration merging 扩展 HookContextMap,这样 ctx.middleware()ctx.runHook() 都能获得精确类型。