news 2026/9/16 19:44:24

Electric Agents Webhook Sources 实战:让 Agent 订阅外部事件流并被精准唤醒

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Electric Agents Webhook Sources 实战:让 Agent 订阅外部事件流并被精准唤醒

Electric Agents Webhook Sources 实战:让 Agent 订阅外部事件流并被精准唤醒

【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric

Webhook sources 是 Electric Agents 平台中连接"外部世界"与 Agent 实体的订阅机制:它允许 Agent(如 Horton 运行时)发现并订阅 GitHub、Stripe、邮件、CI 等外部 Webhook 集成产生的事件流,订阅关系持久化在实体 manifest 中,当匹配的外部事件到达时,平台会携带已水合(hydrated)的事件数据唤醒对应实体。读完本文,你将掌握 Webhook source 的契约结构(contract/bucket/filter)、四个内置工具(list_webhook_sources等)的用法、编程式订阅的客户端 API,以及唤醒载荷(wake payload)的完整数据形态,并能将其接入自定义运行时。

整体机制:订阅—持久化—唤醒

Webhook sources 的核心意图是让 Agent 订阅外部事件源,如 GitHub、Stripe、email、CI 或其他 Webhook 集成。整个流程可以概括为三步:

  1. 发现:Agent 调用list_webhook_sources,获取当前可见的 webhook source 契约(contract),其中声明了可订阅的webhookKey、bucket(路径模板分桶)、参数 schema 与可选的命名过滤器;
  2. 订阅:Agent(或宿主代码)调用subscribe_webhook_source,订阅关系写入实体 manifest(kind: "source"sourceType: "webhook"的条目),随实体流持久化;
  3. 唤醒:当被订阅的源产生匹配事件时,实体被唤醒,唤醒载荷(wake payload)中直接携带水合后的 webhook 事件,Horton 会把这些数据放进 trigger message,让模型无需二次查询即可做出反应。

内置的 Horton 运行时默认通过ctx.electricTools暴露 webhook-source 工具,无需额外接线。

契约:WebhookSourceContract 与 Bucket

一个 webhook source 契约描述了 Agent 可以订阅什么。完整类型定义见 webhook-sources.ts:

type WebhookSourceContract = { serviceId?: string webhookKey: string sourceType: "webhook" endpointKey: string status: "active" | "disabled" | "revoked" label: string description?: string agentVisible: boolean buckets: WebhookSourceBucket[] updatedAt?: string revision: number }

各字段的作用:

  • webhookKey:Agent 发起订阅时使用的源标识(如"github");
  • endpointKey/sourceType:底层端点标识与源类型(固定为"webhook");
  • status:仅active状态的源可被订阅——从源码看,resolveWebhookSourceSubscription 会对非agentVisible或非active的契约直接抛出"is not active"错误;
  • agentVisible:控制该契约是否对 Agent 的list_webhook_sources可见;
  • revision:契约版本号,订阅时会以contractRevision记录,便于后续追踪契约变更;
  • buckets:路径模板分桶,见下文。

Bucket 描述路径模板与参数,类型定义见 webhook-sources.ts:

type WebhookSourceBucket = { key: string label: string description?: string pathTemplate: string paramsSchema: Record<string, unknown> eventTypes?: string[] filters?: WebhookSourceFilter[] }

Bucket 的关键机制:

  • pathTemplate使用:name形式的模板占位符。运行时通过 renderWebhookSourceBucketPath 渲染:正则/:([A-Za-z_][A-Za-z0-9_]*)/g逐个替换占位符,占位符值缺失时抛出Missing bucket parameter错误,值会被encodeURIComponent编码;渲染结果不允许出现//、不允许以/开头,也不能为空字符串;
  • paramsSchema是一份 JSON Schema,订阅时用 Ajv 编译并对提交的params做严格校验(见 validateBucketParams),校验器按 schema 对象缓存在WeakMap中,失败时错误信息会带具体路径(如/repo is required);
  • eventTypes可选地描述该桶关注的事件类型;
  • filters是该桶下可用的命名过滤器。

Agent 的标准工作方式是:先调用list_webhook_sources,然后使用其中声明的webhookKeybucketKeyparamsSchema以及可选的filterKey来发起订阅——所有可订阅的表面都以契约声明为准。

内置工具:createWebhookSourceTools

运行时工具工厂可添加四个工具(实现位于 tools/webhook-sources.ts):

工具用途
list_webhook_sources列出实体可订阅的外部 webhook 源。
list_webhook_source_subscriptions列出本实体当前生效的订阅。
subscribe_webhook_source让实体订阅某个源或桶。
unsubscribe_webhook_source按 id 移除一个订阅。

几个实现层面的细节值得注意:

  • 列表来源list_webhook_sources实际调用运行时的listWebhookSources(),其底层是 HTTPGET /_electric/webhook-sources(见 runtime-server-client.ts),返回服务端维护的契约数组;
  • 订阅列表来源list_webhook_source_subscriptions不走网络,而是直接扫描实体本地manifests集合,通过 getWebhookSourceSubscriptions 过滤出kind === "source"sourceType === "webhook"的 manifest 条目并按 id 排序。也就是说,订阅清单是随实体流同步的,跨唤醒(across wakes)依然可用;
  • 幂等等待subscribe_webhook_sourceunsubscribe_webhook_source执行后都会await db.utils.awaitTxId(txid, 10_000)等待本实体流确认该事务,再返回订阅/删除结果,保证工具返回值与本地状态一致;
  • 日志:每个工具调用都经withWebhookSourceToolLogging包装,记录 start / success / failed 三个阶段及参数,便于排查 Agent 行为。

Horton 从内置运行时直接获得这些工具。自定义运行时可以用createWebhookSourceTools()提供它们,或者通过createRuntimeHandler()传入createElectricTools

import { createWebhookSourceTools } from "@electric-ax/agents-runtime/tools" const runtime = createRuntimeHandler({ baseUrl: "http://localhost:4437", registry, createElectricTools: (context) => createWebhookSourceTools(context), })

注意:内置运行时默认还会添加 schedule(定时任务)工具。如果你替换了createElectricTools,想让 Horton 同时保留两种能力时,需要把两套工具都包含进去。createElectricTools的接线点在 process-wake.ts,运行时在处理唤醒时按需构建工具集。

从工具发起订阅:参数、确定性 ID 与生命周期

subscribe_webhook_source接受如下输入(类型见 webhook-sources.ts):

type WebhookSourceSubscriptionInput = { id?: string webhookKey: string bucketKey?: string params?: Record<string, unknown> filterKey?: string lifetime?: SubscriptionLifetime reason?: string }
  • id省略时的确定性派生:运行时调用 buildWebhookSourceSubscriptionId,用webhookKeybucketKey(缺省用"root")、filterKey拼接出规范化前缀(小写化、非法字符替换为-、截断到 80 字符),再对{webhookKey, bucketKey, params, filterKey}的稳定 JSON(键排序序列化)做 FNV-1a 哈希生成后缀。同一组参数重复订阅会派生出同一 id,天然幂等;
  • 省略bucketKey即订阅源根流(root stream),对应工具参数描述中"Omit to subscribe to the source root stream";
  • filterKey只能选择该源/桶声明过的命名过滤器,resolveWebhookSourceSubscription 会校验 filter 是否存在于对应 bucket 的filters列表中;
  • reason是面向人的订阅理由,会随 manifest 持久化,并出现在后续唤醒载荷中。

生命周期(lifetime)有三种取值:

type SubscriptionLifetime = | { kind: "until_entity_stopped" } | { kind: "expires_at"; at: string } // at 为 ISO-8601 绝对时间 | { kind: "manual" }

默认生命周期是until_entity_stopped——订阅随实体存活,实体停止即失效;expires_at允许设置明确到期时间;manual表示需要显式取消。在工具侧,lifetime 用 TypeBox schema 描述(tools/webhook-sources.ts),expires_at.at要求 ISO-8601 字符串。

编程式订阅:createRuntimeServerClient

宿主代码可以不经 Agent,直接通过createRuntimeServerClient()返回的客户端订阅,完整示例:

await client.subscribeToWebhookSource({ entityUrl: "/horton/onboarding", webhookKey: "github", bucketKey: "repo", params: { repo: "electric-sql/electric" }, reason: "Watch repo activity for this session", }) await client.unsubscribeFromWebhookSource({ entityUrl: "/horton/onboarding", id: "github-main", })

listWebhookSources()检查当前可用的契约:

const sources = await client.listWebhookSources()

从客户端实现看(runtime-server-client.ts),这三个方法对应的服务端路由为:

  • GET /_electric/webhook-sources—— 列出契约;
  • PUT <entityRpcPath>/webhook-source-subscriptions/<id>—— 创建/更新订阅,请求体携带webhookKeybucketKeyparamsfilterKeylifetimereason,返回{ txid, subscription }
  • DELETE <entityRpcPath>/webhook-source-subscriptions/<id>—— 按 id 删除订阅,返回{ txid }

客户端在id缺省时同样调用buildWebhookSourceSubscriptionId派生确定性 id,因此工具路径与编程式路径产生的 id 规则一致。服务端的订阅路由行为有对应测试覆盖:webhook-source-subscriptions-route.test.ts。

唤醒载荷:HydratedWebhookSourceWake

当被订阅的源触发时,实体会被唤醒,并携带水合后的 webhook-source 载荷:

type HydratedWebhookSourceWake = { type: "webhook_source_wake" source: string sourceType: "webhook" endpointKey: string webhookKey: string subscription: { id: string bucketKey?: string params: Record<string, unknown> filterKey?: string reason?: string } bucket: string | null changes: Array<{ collection: string kind: "insert" | "update" | "delete" key: string }> events: WebhookEventRow[] missingEventKeys?: string[] }

这个结构是如何被构建出来的?从源码看(process-wake.ts):

  1. webhookSourceWakeInfoFromManifests检查当前唤醒事件:它必须是一个wake事件,且changes中包含webhook_event集合的变更;然后遍历实体 manifests,找到streamUrl与唤醒源一致的 webhook manifest,还原出订阅信息(sourceUrl、endpointKey、webhookKey、subscriptionId、params 等);
  2. 运行时通过wiringConfig.createSourceDb(sourceStreamUrl, ...)对该源流做一次预加载读取,取出events行的实际数据;
  3. buildHydratedWebhookSourceWake 把唤醒中声明的webhook_event变更 key 与实际读到的事件行做匹配,产出events数组;若某些声明的 key 在实际数据中读不到(例如已被清理),会收集到missingEventKeys中,供处理方感知数据缺口。

Handler 可以检查wake.payload,或直接使用常规 agent context。Horton 会把水合后的 webhook-source 数据放入 trigger message 中——具体实现是 context-factory.ts 在构造触发消息时,若存在hydratedWebhookSourceWake,就将其序列化进消息文本,模型可以立即基于事件内容做出反应,而不需要再做第二次查询。相关行为在 process-wake.test.ts 与 context-factory.test.ts 中有测试覆盖。

Manifest 条目:订阅的持久化形态

订阅以manifest行的形式存储,kind: "source",并使用稳定的 manifest key:

webhook-source:<subscription-id>

buildWebhookSourceManifestEntry 展示了完整条目结构:

  • key:即webhook-source:<id>
  • sourceRef<endpointKey>/<bucketPath>(有桶时)或<endpointKey>(根流);
  • config.streamUrl:订阅对应的实际流 URL;
  • config.webhookSource:订阅全貌,含idwebhookKeybucketKeyparamsfilterKeyfilterAppliedcontractRevisionlifetimereasoncreatedBycreatedAt
  • wake:唤醒规则{ on: "change", collections: ["webhook_event"], ops: ["insert"] }——即当webhook_event集合出现 insert 时触发唤醒。

因为订阅是 manifest 条目,它随实体流同步与持久化,这使得实体可以跨唤醒列出和管理自己的订阅(list_webhook_source_subscriptions正是直接读 manifests 实现的)。

过滤器:当前版本的 advisory 语义

filterKey用于选择源声明的命名过滤器,过滤器用于收窄外部 webhook 流。契约层面过滤器条件(WebhookSourceFilterCondition)已支持按collectionsops(insert/update/delete)以及 CEL 表达式where描述。

但需要注意当前的限制:在本版本中,过滤器是 advisory(建议性)的,直到服务端 webhook 过滤器启用为止。订阅成功后条目中记录的filterApplied字段为false(见 resolveWebhookSourceSubscription),工具描述中也明确提示"filters are advisory until server-side webhook filters are enabled"。因此实践建议是:Agent 在处理器中仍应防御性地处理不符合过滤预期的事件(例如在 handler 里自行判断事件类型后再行动),不能假设过滤已在链路上生效。

小结

Webhook sources 用"契约 + 订阅 + manifest 持久化 + 水合唤醒"四个环节,把外部事件源纳入了 Electric Agents 的同步体系:

  • 契约(WebhookSourceContract/WebhookSourceBucket)声明可订阅面,参数经 JSON Schema 校验、桶路径模板严格渲染;
  • 工具与客户端两条订阅路径共用同一套确定性 id 派生与生命周期语义,订阅落盘为webhook-source:<id>manifest 条目;
  • 唤醒时运行时自动预加载源流、匹配事件行并构建HydratedWebhookSourceWake,Horton 将其注入 trigger message,实现"事件到达即上下文就绪"。

如需进一步跟进实现细节,可参考 webhook-sources.ts、tools/webhook-sources.ts、runtime-server-client.ts,以及测试 webhook-sources.test.ts、webhook-source-tools.test.ts。

【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/16 19:43:32

MTK DRM显示驱动初始化:从Component机制到KMS对象创建全解析

第一次碰MTK平台的DRM显示驱动&#xff0c;大多数人都是在probe流程里迷路的。顶层明明挂着标准drm_driver的皮&#xff0c;真正干活的却是一套叫component的机制——匹配、绑定、拆解绕一大圈&#xff0c;最后才回到KMS对象的创建。再加上MTK显示管线本身又是OVL、RDMA、COLOR…

作者头像 李华
网站建设 2026/9/16 19:42:14

国产Linux系统上用pyenv实现Python多版本隔离管理实战

先说一个我自己的踩坑经历。去年在一台统信UOS 20机器上部署数据分析服务&#xff0c;项目依赖锁文件要求Python 3.10&#xff0c;机器自带的是Python 3.7。当时图省事&#xff0c;直接改了系统Python的软链接&#xff0c;结果重启后桌面环境直接起不来了。后来才搞清楚&#x…

作者头像 李华
网站建设 2026/9/16 19:41:32

智能算法在栅格地图路径规划中的对比与应用

1. 项目背景与核心价值在机器人导航、物流配送和自动驾驶等领域&#xff0c;路径规划始终是核心问题之一。二维栅格地图作为最常见的环境建模方式&#xff0c;其路径优化效果直接影响系统性能。传统算法如A*、Dijkstra在复杂环境中容易陷入局部最优&#xff0c;而智能优化算法因…

作者头像 李华