首页
/ LobeHub Agent Signal 实战:编写 Source、Signal、Action 处理器与策略组合

LobeHub Agent Signal 实战:编写 Source、Signal、Action 处理器与策略组合

2026-09-06 11:39:43作者:何举烈Damon

在 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 —— 因果链元数据(chainIdparentNodeIdrootSourceId 等),让每个 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_totalagent_signal_handler_duration_ms 等指标(见 dispatchHandlers 的文档注释),因此你编写的每个处理器天然具备可观测性。

Fluent 注册 API:四个定义函数

处理器注册统一使用 runtime/middleware.ts 提供的四个 helper:

  • defineSourceHandler(...) —— 注册 source 处理器,监听原始事件;
  • defineSignalHandler(...) —— 注册 signal 处理器,监听语义信号;
  • defineActionHandler(...) —— 注册 action 处理器,监听待执行动作;
  • defineAgentSignalHandlers(...) —— 把一组处理器打包成一个 AgentSignalMiddleware(策略)。

它们完成两件事:

  1. 让注册保持简短:返回值是一个纯数据定义 { type, listen, id, handle }(见 AgentSignalSourceHandlerDefinition),不直接触碰注册表;
  2. 保持强类型:当 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(含 lastEventAtstartedAt,见 context.ts),用于"距上次事件多久"一类的节流判断;
  • runtimeState.touchGuardState(lane, now?):写入/更新守卫状态,now 缺省时使用 ctx.now()

从源码结构看,guard state 的持久化后端是可插拔的:AgentSignalRuntime 内置了内存版 createInMemoryRuntimeGuardBackendAgentSignalRuntime.ts),生产环境则使用 redisGuard.ts 的 Redis 实现,两者都实现同一份 RuntimeGuardBackend 接口。

返回契约:六种返回值如何驱动链条

handler 的返回值类型是 AgentSignalSchedulerHandlerResultAgentSignalScheduler.ts),即 ExecutorResult | RuntimeProcessorResultRuntimeProcessorResultpackages/agent-signal/src/base/types.ts 中定义为四种状态之一的联合,加上"不返回任何值"与 action 处理器专用的 ExecutorResult,完整契约如下:

返回值 语义 调度器行为
void 不产生扇出 链条在本 handler 处停止扩展(不入队新节点)
{ status: 'dispatch', signals?, actions? } 继续链条 signalsactions 全部推入队列继续处理
{ status: 'wait', pending? } 暂停 等待后续宿主(host)协调,返回 pending 状态
{ status: 'schedule', nextHop } 排程 调度另一跳(nextHop 为下一跳描述)
{ status: 'conclude', concluded? } 终结 以终态运行时结果停止
ExecutorResult 仅 action 处理器 报告"已执行具体副作用"的结构化结果

调度器对每种返回值的处理可以逐一对应在 applyResult

  • dispatchresult.signalsresult.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: ExecutorErrorcodemessage、可选 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)] : []),
]);

这个组合里有几个值得学习的写法:

  1. 条件装配:可选依赖(options.procedureoptions.skillManagementoptions.userMemory)用展开运算条件地加入数组,没有对应依赖就不安装该 handler——"缺失可选选项意味着对应 handler 不安装";
  2. 分层清晰:数组里同时出现 source 处理器(createFeedbackSatisfactionJudgeProcessor)、signal 处理器(createFeedbackDomainJudgeSignalHandlercreateFeedbackActionPlannerSignalHandler)和 action 处理器(defineUserMemoryActionHandler),它们共同构成一条"反馈 -> 满意度判断 -> 领域路由 -> 动作规划 -> 执行"的完整链条;
  3. 依赖注入:分类器、数据库、诊断通道都通过 options 注入,便于测试时替换为 stub。

组合好的策略最终通过策略栈工厂进入运行时,见 policies/index.ts

export const createDefaultAgentSignalPolicies = (
  options: CreateDefaultAgentSignalPoliciesOptions = {},
): AgentSignalMiddleware[] => {
  return DEFAULT_AGENT_SIGNAL_POLICY_FACTORIES.flatMap((createPolicy) => createPolicy(options));
};

默认的三条策略工厂分别是 createAnalyzeIntentPolicycreateReviewNightlyPolicycreateCompletionPolicy,返回一个 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.tsfeedbackAction.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.tsAGENT_SIGNAL_POLICY_SIGNAL_TYPES 中。

signal handler 的适用场景:路由(routing)、扇出(fan-out)、过滤(filtering)、冲突消解(conflict resolution)、以及把解释结果转化为规划好的动作。

Action Handler 模式:执行真实副作用

当运行时准备执行具体副作用(或入队一个带外执行的任务)时,写 action handler。参考实现是 actions/userMemory.tsactions/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 类型结构上得到印证):

  1. 幂等检查放在这里(或副作用发生前)。仓库提供了现成的工具 policies/actionIdempotency.tshasAppliedActionIdempotency / markAppliedActionIdempotency),userMemory.ts 就是直接引用它做已执行检测的;
  2. 返回稳定的 actionId:直接使用 action.actionId,它是结果信号(signalId 形如 ${actionId}:result,见 buildExecutionSignal)与幂等键的锚点;
  3. 失败细节写进 error:失败时返回 { status: 'failed', error: { code, message, cause? } }userMemory.ts 中的 toExecutorError 是标准写法(固定错误码 USER_MEMORY_EXECUTION_FAILED + 完整 cause);
  4. ExecutorResult 交给调度器:由 AgentSignalScheduler 把它投影成 actionApplied / actionSkipped / actionFailed 内置结果信号,而不是自己再发一遍;
  5. 异步 execAgent 类动作:handler 内只报告入队(enqueue)结果,持久的执行回执(receipt)从后续 agent.execution.completed source 事件投影,不要在 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.tspackages/agent-signal/src/source/scopeKey.ts
服务端 source 归一化与 hydration apps/server/src/services/agentSignal/sources/buildSource.tshydration/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_TYPESaction.nudge.handleaction.persona.handleaction.skill-management.handleaction.user-memory.handle),每个 action 类型都有对应的 payload 契约类型,defineActionHandler 的输入类型收窄正是依赖这份映射。

如何选择节点类型

规范给出的判断标准:

  • source:外部世界产生了新事实(用户消息、执行完成、工具结果……);
  • signal:系统需要一个可被下游复用的语义含义("这条反馈不满意"、"该走记忆域");
  • action:运行时已准备好执行一个具体副作用。

补充规则:如果一个 handler 既做解释又做副作用,拆成两个——先 source/signal 处理器产出语义,再 action 处理器执行副作用。analyzeIntent 策略正是这一点的范例:satisfaction judge(解释)与 user-memory / skill-management handler(副作用)分离,使链条可检视、可测试。

测试策略

优先在改动代码旁边写聚焦测试。仓库内可直接参考的测试位置:

分层原则是"在能证明行为的最小层级写测试":

  • 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 三条策略就是最完整的可读样例。

登录后查看全文
热门项目推荐
相关项目推荐