LobeHub Agent Signal 实战:编写 Source、Signal、Action 处理器与策略组合
在 LobeHub 的 Agent Signal 运行时中,处理器(Handler)是语义事件链上的最小执行单元:它接收一个运行时节点,返回一个决定链条继续或终止的结果。本文基于仓库内的处理器编写规范与真实实现,讲解如何用 defineSourceHandler / defineSignalHandler / defineActionHandler / defineAgentSignalHandlers 这组 fluent 注册 API 编写类型安全的处理器,如何用策略组合(policy composition)把它们装配进运行时,以及三类处理器各自的编写模式、返回契约、类型归属位置和测试策略——读完后你可以在 apps/server/src/services/agentSignal 下独立编写、注册和验证一个新的语义处理链。
处理器运行在什么样的运行时上
写处理器之前,先理解它面对的运行时形态。LobeHub 的 Agent Signal 运行时把三种节点统一抽象为 RuntimeNode,定义在 runtime/context.ts:
export type RuntimeNode = AgentSignalSource | BaseAction | BaseSignal;
三者都继承自 packages/agent-signal/src/base/types.ts 中的基础契约,每个节点都携带:
chain: AgentSignalChainRef—— 因果链元数据(chainId、parentNodeId、rootSourceId等),让每个 signal/action 都能回溯到自己的来源 source;scopeKey/scope—— 作用域标识,运行时据此做队列隔离与守卫(guard state);timestamp—— 节点产生时间。
调度由 AgentSignalScheduler 完成:它维护一个节点队列,按 sourceType / signalType / actionType 在对应注册表中匹配处理器并逐个执行(见 dispatchNode)。处理器的返回值会被 applyResult 解释为"继续入队"或"终止运行",这正是下一节要讲的返回契约。
每个 handler 执行时还会被包进一个名为 agent_signal.handler.run 的 OpenTelemetry span,并记录 agent_signal_handler_runs_total、agent_signal_handler_duration_ms 等指标(见 dispatchHandlers 的文档注释),因此你编写的每个处理器天然具备可观测性。
Fluent 注册 API:四个定义函数
处理器注册统一使用 runtime/middleware.ts 提供的四个 helper:
defineSourceHandler(...)—— 注册 source 处理器,监听原始事件;defineSignalHandler(...)—— 注册 signal 处理器,监听语义信号;defineActionHandler(...)—— 注册 action 处理器,监听待执行动作;defineAgentSignalHandlers(...)—— 把一组处理器打包成一个AgentSignalMiddleware(策略)。
它们完成两件事:
- 让注册保持简短:返回值是一个纯数据定义
{ type, listen, id, handle }(见 AgentSignalSourceHandlerDefinition),不直接触碰注册表; - 保持强类型:当
listen指向具体的 source / signal / action 类型常量时,handler 的输入参数会被自动收窄到对应的具体 payload 类型。
类型收窄机制值得细看。以 source 为例,middleware.ts 中的条件类型链是:
type ResolveSourceHandlerInput<TListen extends OneOrMany<string>> =
ExtractListenType<TListen> extends infer TType extends string
? TType extends AgentSignalSourceType
? AgentSignalSourceVariant<TType> // 具体 source 类型 -> 具体 payload
: AgentSignalSource // 非内置类型 -> 通用基础类型
: never;
也就是说,当你传入 AGENT_SIGNAL_SOURCE_TYPES.agentUserMessage 时,handler 回调收到的 source 已经是带有 agentUserMessage 特定 payload 结构的类型;传入的 listen 也可以是字符串数组,表示同一个处理器同时监听多个类型。注册发生时,AgentSignalMiddleware.install(context) 会把定义按类型分发给 handleSource / handleSignal / handleAction,而安装上下文 createAgentSignalMiddlewareInstallContext 会把 listen 展开成数组,对每个类型调用对应注册表的 register——一个处理器监听 N 个类型,就会在注册表里登记 N 次。
listen 之外还需要一个 handler id:它既是注册表的登记键,也是可观测性 span 上的 agent.signal.handler_id 属性,命名惯例是 "<node-type>:<语义后缀>",例如真实实现中的 `${AGENT_SIGNAL_SOURCE_TYPES.agentUserMessage}:feedback-satisfaction-judge`(见 feedbackSatisfaction.ts)。
RuntimeProcessorContext:处理器拿到什么
每个 handler 的回调统一接收两个参数:当前运行时节点和 RuntimeProcessorContext。后者的完整定义在 context.ts:
export interface RuntimeProcessorContext {
now: () => number;
runtimeState: RuntimeScopedState;
scopeKey: string;
}
export interface RuntimeScopedState {
getGuardState: (lane: string) => Promise<RuntimeGuardState>;
touchGuardState: (lane: string, now?: number) => Promise<RuntimeGuardState>;
}
逐项说明:
scopeKey:当前链的作用域键,用于跨 handler 定位同一 scope 的持久化状态;now():时间函数。注意调度器在 dispatchHandlers 中把now固定为节点自身的timestamp,即同一条链内所有 handler 看到的是同一个时间基准,保证基于时间的判断(去抖、窗口)可复现;runtimeState.getGuardState(lane):读取 scope 内某条"守卫通道"(guard lane)的状态,返回RuntimeGuardState(含lastEventAt与startedAt,见 context.ts),用于"距上次事件多久"一类的节流判断;runtimeState.touchGuardState(lane, now?):写入/更新守卫状态,now缺省时使用ctx.now()。
从源码结构看,guard state 的持久化后端是可插拔的:AgentSignalRuntime 内置了内存版 createInMemoryRuntimeGuardBackend(AgentSignalRuntime.ts),生产环境则使用 redisGuard.ts 的 Redis 实现,两者都实现同一份 RuntimeGuardBackend 接口。
返回契约:六种返回值如何驱动链条
handler 的返回值类型是 AgentSignalSchedulerHandlerResult(AgentSignalScheduler.ts),即 ExecutorResult | RuntimeProcessorResult。RuntimeProcessorResult 在 packages/agent-signal/src/base/types.ts 中定义为四种状态之一的联合,加上"不返回任何值"与 action 处理器专用的 ExecutorResult,完整契约如下:
| 返回值 | 语义 | 调度器行为 |
|---|---|---|
void |
不产生扇出 | 链条在本 handler 处停止扩展(不入队新节点) |
{ status: 'dispatch', signals?, actions? } |
继续链条 | signals 与 actions 全部推入队列继续处理 |
{ status: 'wait', pending? } |
暂停 | 等待后续宿主(host)协调,返回 pending 状态 |
{ status: 'schedule', nextHop } |
排程 | 调度另一跳(nextHop 为下一跳描述) |
{ status: 'conclude', concluded? } |
终结 | 以终态运行时结果停止 |
ExecutorResult |
仅 action 处理器 | 报告"已执行具体副作用"的结构化结果 |
调度器对每种返回值的处理可以逐一对应在 applyResult:
dispatch:result.signals和result.actions依次queue.push(...),链条继续;conclude/schedule/wait:先记录agent_signal_terminal_results_total终端指标(recordTerminalResultMetric),再把该结果作为本次运行的终端结果返回;ExecutorResult(通过 isExecutorResult 识别applied | failed | skipped状态):被记入 trace,并由调度器自动构造一个执行结果信号推入队列——buildExecutionSignal 会把ExecutorResult转成actionApplied/actionSkipped/actionFailed三种内置结果信号。
最后一点是编写 action 处理器时的关键约定:不要让 action 处理器自己再发一次"我成功了"的信号,把 ExecutorResult 交还给调度器,它会负责投影为内置结果信号,下游 handler 统一消费 AGENT_SIGNAL_TYPES.actionApplied 等即可。
ExecutorResult 本身的结构也在 base/types.ts 中固化:成功/跳过时是 { status: 'applied' | 'skipped', actionId, attempt, output?, detail?, emittedSignalIds? },失败时必须携带结构化 error: ExecutorError(code、message、可选 retriable / retryAfterMs / cause)。attempt 记录本次执行尝试的 startedAt / completedAt / current 次数与状态,是重试与幂等追踪的基础。
策略组合模式:defineAgentSignalHandlers
单个处理器只是零件,真实系统里相关处理器被 defineAgentSignalHandlers([...]) 打包成一个"策略"(AgentSignalMiddleware)。最典型的组合是 analyzeIntent 策略,见 policies/analyzeIntent/index.ts:
return defineAgentSignalHandlers([
...(options.procedure ? [createToolOutcomeSourceHandler(options.procedure)] : []),
createFeedbackSatisfactionJudgeProcessor({
...options.feedbackSatisfactionJudge,
classifierDiagnostics: options.feedbackSatisfactionJudge?.classifierDiagnostics ?? options.classifierDiagnostics,
}),
createFeedbackDomainJudgeSignalHandler({
classifierDiagnostics: options.classifierDiagnostics,
resolveDomains: feedbackDomainResolver,
skillIntentClassifier,
}),
createFeedbackActionPlannerSignalHandler({
markerReader: options.procedure?.markerReader,
procedure: options.procedure,
}),
...(options.skillManagement
? [
defineSkillManagementActionHandler({ ... }),
createCompletionSkillSynthesisSourceHandler({ ... }),
]
: []),
...(options.userMemory ? [defineUserMemoryActionHandler(options.userMemory)] : []),
]);
这个组合里有几个值得学习的写法:
- 条件装配:可选依赖(
options.procedure、options.skillManagement、options.userMemory)用展开运算条件地加入数组,没有对应依赖就不安装该 handler——"缺失可选选项意味着对应 handler 不安装"; - 分层清晰:数组里同时出现 source 处理器(
createFeedbackSatisfactionJudgeProcessor)、signal 处理器(createFeedbackDomainJudgeSignalHandler、createFeedbackActionPlannerSignalHandler)和 action 处理器(defineUserMemoryActionHandler),它们共同构成一条"反馈 -> 满意度判断 -> 领域路由 -> 动作规划 -> 执行"的完整链条; - 依赖注入:分类器、数据库、诊断通道都通过 options 注入,便于测试时替换为 stub。
组合好的策略最终通过策略栈工厂进入运行时,见 policies/index.ts:
export const createDefaultAgentSignalPolicies = (
options: CreateDefaultAgentSignalPoliciesOptions = {},
): AgentSignalMiddleware[] => {
return DEFAULT_AGENT_SIGNAL_POLICY_FACTORIES.flatMap((createPolicy) => createPolicy(options));
};
默认的三条策略工厂分别是 createAnalyzeIntentPolicy、createReviewNightlyPolicy 和 createCompletionPolicy,返回一个 AgentSignalMiddleware[] 数组,随后由 createAgentSignalRuntime({ policies }) 统一安装(运行时入口与工厂位于 runtime/AgentSignalRuntime.ts,集成装配见 orchestrator.ts)。
Source Handler 模式:把原始事件解释成语义信号
当你需要把一个生产者事件(producer event)解释成语义信号时,写 source handler。参考实现是 feedbackSatisfaction.ts:
return defineSourceHandler(
AGENT_SIGNAL_SOURCE_TYPES.agentUserMessage,
`${AGENT_SIGNAL_SOURCE_TYPES.agentUserMessage}:feedback-satisfaction-judge`,
async (source, ctx): Promise<RuntimeProcessorResult | void> => {
// 1. 用模型 judge 对 source payload 做满意度分类(结构化输出 + Zod 校验)
const classifier: SatisfactionClassifierService = {
async classify(input) {
const payload = await judge.judgeSatisfaction(input);
return FeedbackSatisfactionStagePayloadSchema.parse(payload);
},
};
// 2. 分类结果 -> 语义信号;无需继续时返回 void 终止
const result = await classifySatisfaction(source, ctx, {
diagnostics: options.classifierDiagnostics,
satisfactionClassifier: classifier,
});
if (result.type === 'continue') {
return transitionToSignals(result.value).result;
}
return result.result;
},
);
这个真实例子的流程是:用户消息 agent.user.message -> 满意度 judge(模型结构化输出经 Zod 校验)-> 产出 signal.feedback.satisfaction 信号(dispatch);判定为无需处理时直接返回 void 停止。其文件头部的文档注释(feedbackSatisfaction.ts)明确标注了 upstream(agent.user.message)与 downstream(signal.feedback.satisfaction),这是仓库内 handler 的通行文档惯例,值得照做。
何时写 source handler:
- 一条原始消息、生命周期事件或 bot 入口事件需要解释;
- 工作仍然停留在语义层面,不产生副作用。
Signal Handler 模式:语义状态的分支与规划
当一个语义状态需要分支成更多语义状态、或规划出待执行动作时,写 signal handler。参考实现是 feedbackDomain.ts 与 feedbackAction.ts。基本形态:
return defineSignalHandler(
MY_SIGNAL_TYPE,
'signal.my-policy-router',
async (signal): Promise<RuntimeProcessorResult | void> => {
return {
actions: [/* planned work */],
status: 'dispatch',
};
},
);
以领域 judge 为例,createFeedbackDomainJudgeSignalHandler 监听 signal.feedback.satisfaction,调用领域分类后按目标扇出(fan-out)为 signal.feedback.domain.memory / signal.feedback.domain.prompt / signal.feedback.domain.skill / signal.feedback.domain.none 之一(见 feedbackDomain.ts 的工作流注释)。这些信号类型常量集中定义在 policies/types.ts 的 AGENT_SIGNAL_POLICY_SIGNAL_TYPES 中。
signal handler 的适用场景:路由(routing)、扇出(fan-out)、过滤(filtering)、冲突消解(conflict resolution)、以及把解释结果转化为规划好的动作。
Action Handler 模式:执行真实副作用
当运行时准备执行具体副作用(或入队一个带外执行的任务)时,写 action handler。参考实现是 actions/userMemory.ts 与 actions/skillManagement.ts。基本形态:
return defineActionHandler(
MY_ACTION_TYPE,
'action.my-policy-executor',
async (action, ctx): Promise<ExecutorResult> => {
const startedAt = Date.now();
// 执行 service / tool / model 副作用
// 必要时先做幂等检查
return {
actionId: action.actionId,
attempt: {
completedAt: ctx.now(),
current: 1,
startedAt,
status: 'succeeded',
},
status: 'applied',
};
},
);
编写 action handler 时必须遵守的规则(这些规则同时能从 ExecutorResult 类型结构上得到印证):
- 幂等检查放在这里(或副作用发生前)。仓库提供了现成的工具 policies/actionIdempotency.ts(
hasAppliedActionIdempotency/markAppliedActionIdempotency),userMemory.ts就是直接引用它做已执行检测的; - 返回稳定的
actionId:直接使用action.actionId,它是结果信号(signalId形如${actionId}:result,见 buildExecutionSignal)与幂等键的锚点; - 失败细节写进
error:失败时返回{ status: 'failed', error: { code, message, cause? } },userMemory.ts中的 toExecutorError 是标准写法(固定错误码USER_MEMORY_EXECUTION_FAILED+ 完整 cause); - 把
ExecutorResult交给调度器:由AgentSignalScheduler把它投影成actionApplied/actionSkipped/actionFailed内置结果信号,而不是自己再发一遍; - 异步
execAgent类动作:handler 内只报告入队(enqueue)结果,持久的执行回执(receipt)从后续agent.execution.completedsource 事件投影,不要在 action handler 里同步等待。
关于记忆与技能自迭代动作,仓库给出了明确约定:具体副作用是 enqueueSelfIterationRun(...) 入队一个后台运行(实现见 services/selfIteration/dispatch/enqueueSelfIterationRun.ts)。后台运行会打上 Agent Signal operation marker、在 agent runtime 中写入持久资源,然后把变更结果暴露在 completion source 的 selfIteration payload 上(投影逻辑在 services/selfIteration/completion/)。不要在 action handler 里再加一份同步回执投影——单一路径避免重复记账。
Source、Signal、Action 类型的归属位置
新增类型时必须放对位置,规范如下:
| 类型类别 | 位置 |
|---|---|
| 外部事件 payload | packages/agent-signal/src/source/sourceTypes.ts |
| source 事件信封与 scope key | packages/agent-signal/src/source/sourceEvent.ts、packages/agent-signal/src/source/scopeKey.ts |
| 服务端 source 归一化与 hydration | apps/server/src/services/agentSignal/sources/(buildSource.ts、hydration/、renderers/) |
| 策略自有的 signal / action payload | apps/server/src/services/agentSignal/policies/types.ts |
| 归一化共享节点契约 | packages/agent-signal/src/base/types.ts |
核心原则:不要把应用专属的 signal 目录塞进 packages/agent-signal,该包必须保持通用、可复用。策略层类型的实际例子是 policies/types.ts 中的 AGENT_SIGNAL_POLICY_ACTION_TYPES(action.nudge.handle、action.persona.handle、action.skill-management.handle、action.user-memory.handle),每个 action 类型都有对应的 payload 契约类型,defineActionHandler 的输入类型收窄正是依赖这份映射。
如何选择节点类型
规范给出的判断标准:
- source:外部世界产生了新事实(用户消息、执行完成、工具结果……);
- signal:系统需要一个可被下游复用的语义含义("这条反馈不满意"、"该走记忆域");
- action:运行时已准备好执行一个具体副作用。
补充规则:如果一个 handler 既做解释又做副作用,拆成两个——先 source/signal 处理器产出语义,再 action 处理器执行副作用。analyzeIntent 策略正是这一点的范例:satisfaction judge(解释)与 user-memory / skill-management handler(副作用)分离,使链条可检视、可测试。
测试策略
优先在改动代码旁边写聚焦测试。仓库内可直接参考的测试位置:
- runtime/tests/AgentSignalRuntime.test.ts —— 运行时队列与 fan-out;
- __tests__/index.integration.test.ts —— 服务端入口与可观测性持久化的集成验证;
- policies/analyzeIntent/__tests__/ —— 各路由规则的 handler 级单测(
feedbackSatisfaction.test.ts、feedbackDomain.test.ts、feedbackAction.test.ts等); - policies/analyzeIntent/actions/__tests__/ —— action 处理器单测(
userMemory.test.ts、skillManagement.test.ts); - services/selfIteration/completion/__test__/ —— 异步 memory/skill 回执的 completion 投影测试。
分层原则是"在能证明行为的最小层级写测试":
- handler 单测:验证一条路由规则(如满意度判定结果如何映射为信号);
- runtime 测试:验证队列 fan-out 与终端结果;
- completion 投影测试:验证异步 memory/skill 回执从
agent.execution.completed的正确投影; - 集成测试:验证服务入口装配与可观测性持久化。
小结
LobeHub 的 Agent Signal 处理器体系可以用三句话概括:用四个 define*Handler helper 写出类型安全的节点处理器,用 defineAgentSignalHandlers 把相关处理器组合成策略并经 createDefaultAgentSignalPolicies + createAgentSignalRuntime 安装进运行时,用严格的返回契约(void / dispatch / wait / schedule / conclude / ExecutorResult)控制链条的推进与终止。掌握这套模式后,仓库中 policies/ 下的 analyzeIntent、reviewNightly、completion 三条策略就是最完整的可读样例。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0624
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00