资讯动态

@electric-ax/agents-runtime 源码级解读:在 Durable Streams 之上构建持久化 Agent 运行时

发布时间:2026/9/15 20:51:01 来源:尧图企业网站定制
electric-ax/agents-runtime 源码级解读在 Durable Streams 之上构建持久化 Agent 运行时【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electricelectric-ax/agents-runtime是 Electric Agents 的核心行为运行时behavioral stack它让每个 Agent 实体entity拥有自己的 append-only 持久流durable stream每当实体被唤醒wake运行时把该流物化为类型化的 TanStack DB 集合构造出HandlerContext再调用实体唯一入口handler(ctx, wake)。本文以 packages/agents-runtime/README.md 为骨架结合仓库源码create-handler.ts、process-wake.ts、entity-stream-db.ts、define-entity.ts 等深入展开帮助你理解实体如何声明、状态如何读写、webhook 如何接入、信号如何控制生命周期并能在自己的服务里独立部署这套运行时。面向场景为什么需要行为运行时仓库根目录 README.md 将本仓库定位为 The agent platform built on sync。与一次性调一次 LLM 就结束的脚本不同Agent 平台的实体需要跨多次唤醒保持状态一次handler调用可能被中断、重启、横向迁移但下一次唤醒必须能接续上下文。agents-runtime的解法是每个实体独占一条 append-only 流所有事件inbox 消息、运行记录、状态变更、信号都追加到这条流上天然具备持久化与可重放能力每次唤醒时运行时把流物化成类型化的 StreamDB/TanStack DB 集合handler 读到的就是一份可查询的内存镜像状态写入以事件形式追加回流由 Durable Streams 保证顺序与幂等。因此当前 API没有独立的setup()或loader()阶段README 明确说明一切初始化逻辑都放在handler内通过ctx.firstWake判断是否首次运行。安装pnpm add electric-ax/agents-runtime根据 package.jsonpeer dependency 为tanstack/db 0.5.33README 声明实际开发依赖tanstack/db ^0.6.6运行时还内置了durable-streams/client/durable-streams/state、anthropic-ai/sdk、mariozechner/pi-agent-core/pi-aiLLM 执行层、zod、xstate、cron-parser等依赖说明这是一个自带模型调用、调度、沙箱能力的完整运行时模块同时导出./react、./client、./tools、./sandbox、./sandbox/docker子路径见 package.json 的exports字段方便按需引入react、tanstack/react-db、dockerode、e2b均为可选peer dependency。快速开始一个可运行的 durable chat assistantREADME 的 Quick Start 给出了完整可运行的骨架直接继承如下import http from node:http import { z } from zod import { createRuntimeHandler, defineEntity } from electric-ax/agents-runtime defineEntity(assistant, { description: Simple durable chat assistant, state: { status: { schema: z.object({ key: z.string(), value: z.enum([idle, working]), }), type: status, primaryKey: key, }, }, async handler(ctx) { if (!ctx.db.collections.status.get(current)) { ctx.db.actions.status_insert({ row: { key: current, value: idle }, }) } ctx.useAgent({ systemPrompt: You are a helpful assistant., model: claude-sonnet-4-5-20250929, tools: [...ctx.electricTools], }) await ctx.agent.run() }, }) const runtime createRuntimeHandler({ baseUrl: http://127.0.0.1:4437, serveEndpoint: http://127.0.0.1:3000/webhook, }) await runtime.registerTypes() const server http.createServer(async (req, res) { if (req.url /webhook req.method POST) { await runtime.onEnter(req, res) return } res.writeHead(404) res.end() }) server.listen(3000, 127.0.0.1)关键流程拆解defineEntity声明实体类型实体名 描述 状态集合 handler。注册逻辑在 define-entity.ts 的EntityRegistry中重名会直接抛错already registeredgetEntityType/listEntityTypes/clearRegistry提供查询与重置。createRuntimeHandler创建 webhook 入口baseUrl指向 Durable Streams 服务器serveEndpoint是你的应用暴露的回调地址。实现位于 create-handler.ts它先创建createRuntimeRouter再包装出 Node HTTP 适配器onEnter。registerTypes()注册实体类型运行时把每个实体类型含 Zod schema 转换成的 JSON Schema、state_schemas、serve_endpoint、默认 dispatch policyPOST 到{baseUrl}/_electric/entity-types见 create-handler.ts。默认并发注册数 8registrationConcurrency任何一个类型注册失败会抛错并报告已成功数量。webhook 分发收到 POST 后校验签名、解析WebhookNotification按entity.type找到注册过的实体类型后派发 wake 并立即返回{ ok: true }200真正的 handler 在后台执行。注意ctx.useAgent中用到的electricTools是运行时内置工具集可在createElectricTools钩子中定制见下节不调用useAgent()的 wake 不会触发 LLM 调用。Entity Model实体定义的三要素README 的pr-reviewer示例展示了实体定义的完整形态这里完整保留并逐项注解import { z } from zod import { defineEntity } from electric-ax/agents-runtime defineEntity(pr-reviewer, { description: Code review agent, // 1. 创建参数实体被 spawn 时传入的 args 校验 creationSchema: z.object({ repo: z.string(), prNumber: z.number().int(), }), // 2. 收件箱消息 schema按消息类型声明 payload 形状 inboxSchemas: { pr-opened: z.object({ diff: z.string(), baseBranch: z.string().optional(), }), }, // 3. 状态集合类型化持久状态 state: { reviewStatus: { schema: z.object({ key: z.string(), status: z.enum([pending, reviewing, done]), }), type: review_status, primaryKey: key, }, }, async handler(ctx, wake) { if (ctx.firstWake !ctx.db.collections.reviewStatus.get(current)) { ctx.db.actions.reviewStatus_insert({ row: { key: current, status: reviewing }, }) } if (wake.type inbox) { ctx.useAgent({ systemPrompt: Review pull requests carefully and concisely., model: claude-sonnet-4-5-20250929, tools: [...ctx.electricTools], }) await ctx.agent.run() ctx.db.actions.reviewStatus_update({ key: current, updater: (draft) { draft.status done }, }) } }, })三个要素说明creationSchemaspawn 时的参数校验。从源码 types.ts 看校验通过后的 args 会以EntityArgsReadonly不可变对象注入ctx.args。inboxSchemas外部消息进入实体收件箱inbox时的 payload 校验handler 里通过wake.type inbox区分消息类型。state声明式状态集合。每个集合包含schema任意 Standard Schema常见用 zod、type事件类型名默认state:${name}、primaryKey主键字段默认key。类型定义见 types.ts该处还支持externallyWritable允许外部通过 HTTP 写集合写者身份会被物化为_principal虚拟列、contract、operations等高级选项。在注册阶段运行时会把所有 Zod schema 通过 create-handler.ts 的toJsonSchema统一转换成 JSON Schema优先读取~standard/toJSONSchema接口兜底用zodToJsonSchema再连同内置的DEFAULT_STATE_SCHEMAS一起作为state_schemas上报给服务器buildEntityTypeRegistrationBody。HandlerContexthandler 的完整工具箱handler(ctx, wake)收到的ctx提供以下能力README 原文要点结合 types.ts 的RuntimeContext补充类型细节ctx.firstWake—— 仅实体首次唤醒时为true用于初始化逻辑ctx.entityUrl/ctx.entityType—— 实体身份元数据ctx.args—— 不可变的、已校验的 spawn 参数ctx.db—— 实体物化后的 StreamDB读用ctx.db.collections.name写用自动生成的ctx.db.actions.name_{insert,update,delete}ctx.electricTools—— 内置运行时工具直接传给 agentctx.spawn/ctx.observe/ctx.send/ctx.mkdb—— 子实体编排与共享状态ctx.self、ctx.attachments、ctx.createEffect等高级句柄。wake参数区分两种类型InboxHandlerWake带message.type和payload的收件箱消息与OtherHandlerWake其他唤醒事件如子实体runFinished见 types.ts。State声明式状态自动生成 CRUD action在definition.state中声明的每个集合都会自动获得ctx.db.actions上的name_insert/name_update/name_delete三个 action读取则通过ctx.db.collections.namectx.db.actions.counts_insert({ row: { key: main, value: 0 } }) ctx.db.actions.counts_update({ key: main, updater: (draft) { draft.value }, }) ctx.db.actions.counts_delete({ key: main }) ctx.db.collections.counts.get(main) ctx.db.collections.counts.toArray这些 action 并非手写而是在 entity-stream-db.ts 中自动生成的onMutate阶段直接改写 TanStack DB 集合乐观更新mutationFn通过共享 producer 把变更以事务事件写回流并等待txid往返确认awaitTxId超时 20 秒见 WRITE_TXID_TIMEOUT_MSupdate对不存在的 key 会退化为insert先以 primaryKey 构造行再执行 updater避免本地乐观事件先于 insert 提交时的顺序缺口。集合定义在创建时与内置集合runs、steps、texts、inbox、manifests、signals、errors 等见 entity-schema.ts合并schema 会包一层虚拟列保活逻辑_timeline_order、_principal确保 TanStack DB 校验时不会丢弃这些运行时附加字段entity-stream-db.ts。Agents模型执行入口ctx.useAgent({ systemPrompt: You are a helpful assistant., model: claude-sonnet-4-5-20250929, tools: [...ctx.electricTools, myTool], }) await ctx.agent.run()agent.run()会基于物化后的时间线组装实体上下文并处理本次唤醒的触发消息。AgentConfig类型见 types.ts支持provider、getApiKey、reasoning、thinkingBudgets、modelTimeoutMs、modelMaxRetries、onStepEnd每步 token 统计回调等配置模型执行由pi-ai提供底层能力。若本次 wake 不调用useAgent()运行时不会触发 LLM。Spawn、Observe 与 Send实体编排原语const child await ctx.spawn( researcher, r-1, { topic: durability }, { wake: { on: runFinished, includeResponse: true } } ) // Return from this wake. Continue when the child completion wake arrives. const observed await ctx.observe(entity(/researcher/r-2), { wake: { on: runFinished, includeResponse: true }, }) const status observed.status() ctx.send(/other-entity/id, { text: hello }, { type: message })spawn(type, id, args, opts)创建子实体返回EntityHandle含entityUrl、db、send、status()等见 types.ts并可注册子实体完成时的唤醒条件runFinished或change支持debounceMs/timeoutMs见Wake联合类型 types.ts。spawn 后先 return等子实体完成唤醒再继续这是典型的长时间运行编排模式。observe(source, { wake })观察任意源实体、共享数据库等返回句柄查询状态。send(entityUrl, payload, opts)向其他实体投递消息type指定消息类型afterMs可延迟投递返回sent或queued。Shared State跨实体的共享状态const boardSchema { findings: { schema: z.object({ key: z.string(), domain: z.string(), finding: z.string(), }), type: finding, primaryKey: key, }, } if (ctx.firstWake) { ctx.mkdb(board-1, boardSchema) } const board await ctx.observe(db(board-1, boardSchema)) board.findings.insert({ key: f-1, domain: security, finding: XSS found, })mkdb(id, schema)创建共享状态流已存在则抛错。SharedStateCollectionSchema与实体state集合的声明方式一致见 types.ts。observe(db(...))返回SharedStateHandletypes.ts提供类型化的集合代理——insert/update/delete返回事务可await tx.isPersisted.promiseget/toArray走底层 TanStack DB 集合。观察者在任意 wake 上都可读写该共享流且 schema 支持type覆盖事件类型名。Runtime Handlerwebhook 入口与配置项createRuntimeHandler()是 Electric Agents 服务器唤醒实体时的 webhook 入口const runtime createRuntimeHandler({ baseUrl: http://127.0.0.1:4437, serveEndpoint: http://127.0.0.1:3000/webhook, }) await runtime.registerTypes()主要方法README 原文附源码位置runtime.registerTypes()—— 向服务器注册全部实体类型create-handler.tsruntime.onEnter(req, res)—— Node HTTP 适配器内部把IncomingMessage转成 fetchRequest后走同一套 webhook 处理逻辑create-handler.tsruntime.handleRequest(request)—— fetch 原生路由非 webhook 路径返回nullcreate-handler.tsruntime.drainWakes()—— 等待所有在途 wake 收敛若有未处理错误会聚合抛出AggregateError配套还有waitForSettled()、abortWakes()、isWakeActive(streamPath)、debugState()create-handler.ts。除baseUrl、serveEndpoint外RuntimeRouterConfigcreate-handler.ts还提供一批值得在生产中使用的配置配置项默认值说明webhookPathserveEndpoint的 pathname兜底/electric-agentshandleRequest匹配的路径handlerUrl—serveEndpoint的向后兼容别名registry模块级默认注册表运行时私有的实体注册表subscriptionPathForType—按类型定制 webhook 订阅路径defaultDispatchPolicyForType回退为{ type: webhook, url: serveEndpoint, ... }类型级默认分发策略serverHeaders—发往 agents-server 控制面请求的附加头webhookSignature启用JWKS 取{baseUrl}/__ds/jwks.json可传false关闭仅限可信测试环境idleTimeout20_000mswake 空闲关闭超时heartbeatInterval10_000ms心跳间隔registrationConcurrency8类型注册并发数createElectricTools内置工具每个 wake 上下文的工具工厂钩子onWakeError—后台 wake 失败观察器返回true视为已处理sandboxProfiles—注册沙箱 profile工厂闭包仅本地持有描述字段上报 UIpublicUrl/namename默认default运行时公共可见性与去重标识webhook 处理内部做了两层防御create-handler.ts先校验webhook-signature头默认从{baseUrl}/__ds/jwks.json拉取 JWKS再校验entity.type已注册JSON 解析失败返回 400未知实体类型返回 503成功则 200 应答并异步执行。Entity Signals用信号控制实体生命周期Agents 可以通过服务器生命周期端点向实体发信号POST /_electric/entities/:type/:instanceId/signal请求体为{ signal: SIGINT, reason: ..., payload: ... }。CLI 暴露相同能力electric agents signal /horton/onboarding SIGINT --reason stop current run electric agents kill /horton/onboardingREADME 给出的信号语义源码侧可在 process-wake.ts 的handleSignalEvent中逐一印证entity-schema.ts 定义了全部信号枚举SIGINT—— 中止当前活动的 handler 调用实体状态保持不变runAbortController.abort() 请求关闭 wake非 Agent handler 可通过ctx.signal观察取消并把它传给可取消任务fetch、子进程SIGSTOP—— 暂停实体新消息仍会唤醒它但暂停期间待处理工作不执行pauseRequested trueSIGCONT—— 恢复被暂停的实体SIGHUP/SIGTERM/SIGUSR—— 在当前 wake 活跃期间投递给ctx.onSignal(...)SIGTERM—— 空闲或暂停实体转为stopped运行中实体进入stopping直至运行时清理完成SIGKILL—— 立即转为killed注销 wakes 并关闭实体流服务器端处理运行时仍会观察并中止在途模型/工具调用。每个信号事件都会被标记为handled并记录outcomeaborted/ignored/shutdown_requested/transitioned/delivered/failed与handled_at写入实体的 signals 集合process-wake.ts便于审计与 UI 展示。内置 Agents 与 Timeline Helpers内置 AgentREADME 明确说明旧的registerChatAgent/registerResearcherAgent/registerCoderAgent/registerOracleAgent辅助函数已移除。内置 AgentHorton worker现在位于electric-ax/agents包在自己的运行时里通过createBuiltinAgentHandler注册。Timeline Helpers运行时额外导出供 UI 层使用的时间线辅助函数createEntityIncludesQuery(db)getEntityState(runs, inbox)buildSections(runs, inbox)timelineToMessages(db)其中useChat(db)使用基于查询的IncludesRun/IncludesInboxMessage数组timelineToMessages(db)从同一形状的数据构建 LLM 消息。这些工具的实现入口在 index.tsbuildSections、buildTimelineEntries、timelineToMessages与 entity-timeline.ts通过./entity-timeline导出createEntityIncludesQuery、getEntityState等见 index.ts。仓库中还提供了 use-chat.ts 与 use-chat-hook.ts 作为 React 侧消费层随./react子路径导出。测试仓库内的集成测试设施README 指出旧的公开createTestAdapterAPI 已移除仓库内的集成测试使用test/runtime-dsl.ts —— 面向 Durable Streams 测试服务器的运行时 DSL 封装StreamHistory提供count/some/find/completedRunCount/snapshot等事件断言工具并加载真实的 agents-server 与测试后端test/runtime-dsl.test.ts —— 基于该 DSL 的端到端行为测试test/setup-context.test.ts —— 上下文装配专项测试。此外 test/ 下还有覆盖状态写入、信号、沙箱、token 预算、上下文装配、webhook 源等维度的 80 个测试文件如entity-stream-db-principal.test.ts、sandbox-docker.test.ts、wake-session.test.ts、webhook-sources.test.ts可作为理解各模块行为的活文档。从源码看一次 wake 的完整旅程把以上散点串起来一次 wake 在 process-wake.ts 的processWake中经历如下阶段实现细节来自源码校验与定位从WebhookNotificationtypes.ts解析entity.type在注册表中查找实体定义未知类型直接忽略且不 ack交给服务器超时回收。打开持久流为实体流创建单个DurableStreamprocess-wake.tsStreamDB、IdempotentProducer幂等写入带epoch与 claim token 自动认领与 SSE tail 共享这条连接。物化状态createEntityStreamDB把自定义state集合与内置集合合并、包 schema、生成 CRUD action通过回放replay将流内历史事件物化为内存集合entity-stream-db.ts。构造上下文createHandlerContext组装ctxdb、spawn/observe/send/mkdb、attachments、electricTools 等处理 inbox 唤醒事件或合并批量 wake 事件wake_batch。执行 handler运行handler(ctx, wake)useAgent后执行 LLM 循环每一步的文本/工具调用/推理事件都写入实体流期间通过idleTimeout20s与heartbeatInterval10s管理长连接生命周期。后台信号与清理信号事件由handleSignalEvent同步处理wake 结束后drainWakes等待全部在途任务宿主关闭时abortWakes快速中止。这条链路也解释了 README 开篇那句话的含义运行时没有独立的 setup/loader 阶段因为流的加载replay 物化是运行时内建动作用户只写handler。小结electric-ax/agents-runtime以Durable Streams TanStack DB 物化 单一 handler 入口三件套为核心defineEntity声明实体创建 schema / 收件箱 schema / 状态集合createRuntimeHandler接入 webhookHandlerContext提供状态读写、Agent 执行、子实体编排与共享状态全套原语Entity Signals提供 Unix 信号式生命周期控制。对想要自建持久化 Agent 服务的开发者而言可以直接用 packages/agents-runtime/README.md 的 Quick Start 起步再按本文给出的源码路径深入 create-handler.ts、process-wake.ts、entity-stream-db.ts 理解每个环节的实现从而在自己的项目中安全地定制工具、沙箱与信号处理逻辑。【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价