EventBus 与 HookRegistry — 事件与中间件
Aalis 提供两种互补的扩展机制:事件(单向通知)和中间件钩子(可拦截的管道)。
EventBus — 事件总线
源码: packages/core/src/primitives/events.ts
类型安全的全局发布/订阅事件总线,用于松耦合的异步通知。事件只通知、不干预流程。
API
// 监听事件(返回 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:starting | — | start() 的第一步 |
app:ready | — | 启动第一相位(sticky) |
app:started | — | 启动第二相位:全部 app:ready 监听器完成之后(sticky);CLI / TUI 在此接管终端 |
app:restarting | — | restart() 先发本事件,监听器全部完成后才把控制交给宿主的 RestartStrategy |
app:stopping | — | stop() 开头,插件拓扑逆序 dispose 之前 |
通知(service:* / plugin:* / plugins:changed):发射方不等监听器。发射点要么是同步的注册 / 拆卸收尾, 要么在 PluginManager 的 recompute flight 或挂起段内——那里等监听器会与 plugins.idle() 互等死锁。 监听器因此不能假设「我返回了状态机才继续」;要看落定后的状态请 await plugins.idle()。
| 事件 | 参数 | 时机 |
|---|---|---|
service:registered | name | 某服务多了一个提供者(ctx.provide) |
service:unregistered | name | 某服务少了一个提供者(退订闭包或 Context 拆卸) |
service:preference-changed | name | 该服务的偏好 provider 切换(preferService / unpreferService);whenService 借此重挂 |
plugin:loaded | instanceId | 插件实例已激活;同一轮 recompute 可能紧接着激活下一个插件 |
plugin:unloaded | instanceId | 插件实例已拆卸;激活失败的回滚与关机拆卸不发 |
plugins:changed | — | 一轮 recompute 收敛,插件状态集合可能已变;关机轮不发 |
app:ready 与 app:started 是两个相位,不是同一里程碑的两个名字:start() 串行 await,app:started 严格晚于 全部 app:ready 监听器完成。app:stopping 用于知会(打印告别语、切状态条),不是清理通道——清理副作用一律走 ctx.onDispose(fn),它覆盖 bounce / unload / 停机全部路径。总线上没有 dispose 事件。
扩展自定义事件
第三方插件通过 TypeScript declaration merging 即可为事件系统新增类型安全的自定义事件:
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+):
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
// 注册中间件(同一钩子内按注册顺序执行)
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-agent | agent:input:before / agent:llm:before / agent:llm:after / agent:tool:before / agent:tool:after / agent:reply:before / agent:turn:after | agent 一轮处理的各相位 |
@aalis/api-gateway | inbound:confirm / inbound:command / inbound:flow / inbound:trigger / inbound:dispatch / outbound:dispatch | 网关出入站的命名相位 |
@aalis/api-memory | memory:clear | 统一记忆清理编排,供 /clear 与各记忆插件协作 |
中间件特性
- 注册顺序执行: 同一钩子内按注册顺序串行执行(无优先级数字;相位间次序由调度方显式表达)
- 数据修改: data 通过引用传递,修改 data 对象即影响后续中间件和 defaultAction
- 流程控制: 调用
next()继续管道;不调用则中止后续中间件和 defaultAction - 上下文绑定: 每个中间件关联注册方 Context 的清理归属,插件卸载时自动清理(
unregisterByOwner)
典型用法
往提示词里加内容不要用本钩子——那是
agent:prompt贡献点的活(ctx.contribute, 见 architecture.md 扩展机制):贡献点自带幂等、确定性排布与 错误隔离,中间件手搓unshift三者全无。本钩子留给改写 / 截停语义。
// 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();
});扩展自定义钩子
第三方插件可以定义自己的钩子,并让其他插件注入中间件:
// 声明类型(可选但推荐)
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() 都能获得精确类型。