
iii 自定义 Trigger Type 开发指南从绑定既有事件源到发布自己的触发器【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii本文是一份面向 iii 工作流worker开发者的触发器Trigger实战指南。你将理解触发器在 iii 中消费与发布两种角色学会用 Node/TypeScript、Python、Rust 三种 SDK 将函数绑定到http、cron、queue、state等既有触发器类型并进一步以发布者身份声明自定义触发器类型Trigger Type、维护绑定路由表、附加 JSON Schema 契约最终通过worker.trigger把事件分发到已绑定的函数。读完本文你能够从零实现一个类似iii-http的迷你 HTTP 触发器发布者并让其他 worker 的函数在你的事件源上运行。理解编写触发器的两种角色一个 worker 使用触发器有两种方式作为消费者Consumer最常见把 worker 自己的函数绑定到系统中已存在的触发器类型上例如http、cron定时执行、queue 消息每条消息触发一次、state状态变化以及其他任何事件源。对应 API 是worker.registerTrigger({ type, function_id, config })。作为发布者Publisher较少见从自己的 worker 注册一个全新的触发器类型让其他 worker 可以把函数绑定到你 worker 发出的事件上例如HTTP 请求到达webhook 命中文件变化数据库更新。本文的主体是第二种如何创造你自己的触发器。如果只是想在新 worker 中使用既有触发器参考 Using iii / Triggers调用侧的机制worker.trigger/iii trigger直接调用、TriggerAction变体、用条件门控、同一函数多个绑定同样见该页。把函数绑定到既有触发器类型大多数 worker 消费其他 worker 已发布的触发器类型http把函数暴露为 HTTP 端点cron让函数按计划执行queue 触发器让函数随每条消息触发state让函数响应数据变化。绑定操作通过worker.registerTrigger({ type, function_id, config })完成Node / TypeScriptworker.registerTrigger({ type: http, function_id: math::add, config: { api_path: /math/add, http_method: POST }, });Pythonworker.register_trigger({ type: http, function_id: math::add, config: {api_path: /math/add, http_method: POST}, })Rustuse iii_sdk::RegisterTriggerInput; use serde_json::json; worker.register_trigger(RegisterTriggerInput { trigger_type: http.into(), function_id: math::add.into(), config: json!({ api_path: /math/add, http_method: POST }), metadata: None, })?;config的具体结构由各触发器类型自行定义并记录在对应发布 worker 的文档中。例如http类型的 config 就是{ api_path, http_method }见后文 schema 一节。一个值得注意的细节是你可以在发布者 worker 尚未连入时就发起注册。引擎会乐观地存储该绑定并在发布者加入网络后自动激活它——这种乐观注册的行为细节同样参见 Using iii / Triggers。其他绑定机制——注销句柄trigger.unregister()、同一函数绑定多个触发器、用condition_function_id做门控、TriggerAction变体Void、Enqueue等——均属于调用侧主题见 Using iii / Triggers。为触发器绑定附加元数据每个触发器绑定都可以携带一个可选的metadataJSON 对象由消费者在注册时设置。引擎原样存储它并在两个地方对外暴露发布者可见发布者的TriggerHandler.registerTrigger(config)回调会以config.metadata收到它发布者可以根据消费者打上的标签做处理——优先级提示、审计标签、路由键等内部记账信息。可被发现engine::triggers::list会在每个TriggerInfo上返回它控制台console以及任何做服务发现的 worker 都能读到。Node / TypeScriptworker.registerTrigger({ type: http, function_id: math::add, config: { api_path: /math/add, http_method: POST }, metadata: { team: platform, env: staging }, });Pythonworker.register_trigger({ type: http, function_id: math::add, config: {api_path: /math/add, http_method: POST}, metadata: {team: platform, env: staging}, })Rustuse iii_sdk::RegisterTriggerInput; use serde_json::json; worker.register_trigger(RegisterTriggerInput { trigger_type: http.into(), function_id: math::add.into(), config: json!({ api_path: /math/add, http_method: POST }), metadata: Some(json!({ team: platform, env: staging })), })?;别混淆metadata 与 schema触发器类型本身没有 metadata 字段——metadata 是按绑定附加的而非按类型。更关键的是不要把它与触发器类型的schematrigger_request_format与call_request_format混为一谈二者的设定方和用途完全不同Metadata由消费者在每次绑定时worker.registerTrigger()调用设置。它是引擎原样存储的自由标签袋服务于发布者的记账与发现。例如绑定到http的消费者可能附带metadata: { team: platform, env: staging, on_call: alice }发布者据此记录团队信息engine::triggers::list也能在请求时返回这些信息。Schemas由发布者在声明触发器类型时设置。它们描述消费者交互的 JSON 形状。例如iii-http发布的http类型会声明config消费者绑定时传入的内容{ api_path, http_method }调用载荷绑定函数在每次请求时收到的东西{ method, headers, query_params, body }。声明一个触发器类型成为发布者前面的内容都是消费者视角你的 worker 函数被绑定到其他 worker 发布的类型上。现在角色反转你的 worker 是发布者你希望其他 worker 注册的函数能在你 worker 观察到的事件HTTP 请求、webhook 命中、文件变化、数据库更新上被触发。触发器类型的组成一个触发器类型由两部分捆绑而成一个字符串id消费者绑定时引用它例如type: mini-http。一个按绑定维护的路由表这个表由你的 worker 在进程内自行维护。引擎的注册表registry会以规范形式记录绑定这正是engine::triggers::list返回的内容但引擎并不会基于它做分发。引擎只负责把网络上任何消费者 worker 的绑定/解绑事件作为回调转发给你的发布者 worker具体如何处置每个绑定由你的 worker 决定。在启动时用worker.registerTriggerType({ id, description }, handler)声明一次触发器类型。你需要实现的TriggerHandler接口暴露两个回调每当有消费者绑定或解绑时引擎会在你的发布者 worker 上调用它们registerTrigger(config)任何消费者 worker 把函数绑定到你的类型时触发。config携带触发器实例的id、消费者的function_id以及符合你类型所接受形状的消费者config。把它存起来。unregisterTrigger(config)解绑时触发从你的表里删除它。触发器类型可以在运行期任意时刻被拆除调用worker.unregisterTriggerType(...)Python 与 Rust 中为worker.unregister_trigger_type(...)签名见下文注销触发器类型一节。示例从零实现一个迷你iii-http下面这个例子勾勒出真实http触发器类型发布者的精简版。发布者 worker 需要做三件事声明一个名为mini-http、形状为 HTTP 的触发器类型维护一张bindings映射表{ trigger id → function_id, methodpath }随着消费者绑定/解绑而增删之后在收到 HTTP 请求时查询正确的绑定并触发。触发绑定函数的细节见下文向已绑定函数分发事件。Node / TypeScriptimport { registerWorker } from iii-sdk; import type { TriggerConfig, TriggerHandler } from iii-sdk; const url process.env.III_URL; if (!url) throw new Error(III_URL must be set); const worker registerWorker(url); type MiniHttpConfig { api_path: string; // leading slash, e.g. /orders http_method?: GET | POST | PUT | DELETE; }; const bindings new Mapstring, TriggerConfigMiniHttpConfig(); const httpHandler: TriggerHandlerMiniHttpConfig { async registerTrigger(config) { bindings.set(config.id, config); }, async unregisterTrigger(config) { bindings.delete(config.id); }, }; worker.registerTriggerType( { id: mini-http, description: Routes HTTP requests to bound functions }, httpHandler, );Pythonimport os from iii import ( InitOptions, RegisterTriggerTypeInput, TriggerConfig, TriggerHandler, register_worker, ) worker register_worker( os.environ.get(III_URL), InitOptions(worker_namemini-http-worker), ) bindings: dict[str, TriggerConfig] {} class HttpHandler(TriggerHandler): async def register_trigger(self, config: TriggerConfig) - None: bindings[config.id] config async def unregister_trigger(self, config: TriggerConfig) - None: bindings.pop(config.id, None) worker.register_trigger_type( RegisterTriggerTypeInput( idmini-http, descriptionRoutes HTTP requests to bound functions, ), HttpHandler(), )Rustuse std::collections::HashMap; use std::sync::{Arc, Mutex}; use iii_sdk::{ InitOptions, RegisterTriggerType, TriggerConfig, TriggerHandler, register_worker, }; let url std::env::var(III_URL).expect(III_URL must be set); let worker register_worker(url, InitOptions::default()); #[derive(Default)] struct HttpHandler { bindings: ArcMutexHashMapString, TriggerConfig, } #[async_trait::async_trait] impl TriggerHandler for HttpHandler { async fn register_trigger(self, config: TriggerConfig) - Result(), iii_sdk::IIIError { self.bindings.lock().unwrap().insert(config.id.clone(), config); Ok(()) } async fn unregister_trigger(self, config: TriggerConfig) - Result(), iii_sdk::IIIError { self.bindings.lock().unwrap().remove(config.id); Ok(()) } } worker.register_trigger_type( RegisterTriggerType::new( mini-http, Routes HTTP requests to bound functions, HttpHandler::default(), ), );从 SDK 源码可以看到TriggerConfig的真实字段除了文档中提到的id、function_id、config、metadata之外还包含一个namespace字段——当注册时省略命名空间SDK 会用注册 worker 自身的命名空间补全发布者若存储了该 config 并在之后调用trigger()必须把这个已解析的命名空间透传下去见 sdk/packages/node/iii/src/triggers.ts、sdk/packages/python/iii/src/iii/triggers.py、sdk/packages/rust/iii/src/triggers.rs。为触发器类型附加 Schema一个触发器类型可以携带两个可选的 JSON Schema用于描述它的载荷trigger_request_format消费者在worker.registerTrigger(...)绑定函数时传入的按绑定config的 schema。call_request_format触发器触发时你的 worker 交付给绑定函数的调用载荷的 schema。两者都会输入 iii 控制台、Agent 可读的 skills以及engine::trigger-types::list的输出让消费者知道该传什么、会收到什么。注意运行时校验目前尚未支持。附加的 schema 仅是信息性的——引擎不会拒绝不符合它们的config值或调用载荷。请把 schema 当作面向消费者、Agent 和控制台的契约文档这与函数请求/响应 schema 的注意事项一致。各 SDK 以自己惯用的方式接收这两个 schemaSDK传入方式Node / Browser在trigger_request_format/call_request_format上传原始 JSON Schema 对象Zod 4 schema 可用z.toJSONSchema(...)转换。Python在RegisterTriggerTypeInput的相同字段上传 Pydantic 模型类自动转换或原始 dict。RustRegisterTriggerType上的构造器方法.trigger_request_format::T()与.call_request_format::T()其中T: schemars::JsonSchema。以引擎内置的http类型为例其真实 schema 定义在 engine/src/trigger_formats.rsHttpTriggerConfig包含api_path如/users/:id、可选的http_method缺省为 GET以及可选的condition_function_idHttpCallRequest则包含query_params、路径参数等字段。Rust SDK 侧RegisterTriggerType构造器会把这些 schema 序列化进注册消息见 sdk/packages/rust/iii/src/iii.rs。注销一个触发器类型当触发器类型所路由的工作不再需要时可以在运行期拆除它。当发布者 worker 断线时它宣告的所有触发器类型都会被自动移除引擎也会停止路由依赖它们的事件——所以显式注销只在worker 保持连接但想丢弃某个类型时才必要。可以在registerTriggerType之后的任意时刻调用前提是 worker 保持连接。典型场景包括底层资源进入维护模式、功能开关关闭了该对外面、或想在不重启的情况下把类型轮换到新 schema。沿用mini-http的例子这里 worker 因其 HTTP 监听器被配置关闭而丢弃mini-httpNode / TypeScript// e.g. config reload disabled the HTTP listener; stop accepting new bindings // while the worker keeps serving other trigger types. worker.unregisterTriggerType({ id: mini-http, description: Routes HTTP requests to bound functions, });Python# e.g. config reload disabled the HTTP listener; stop accepting new bindings # while the worker keeps serving other trigger types. worker.unregister_trigger_type( {id: mini-http, description: Routes HTTP requests to bound functions} )Rust// e.g. config reload disabled the HTTP listener; stop accepting new bindings // while the worker keeps serving other trigger types. worker.unregister_trigger_type(mini-http);三种 SDK 在此处存在签名差异见 sdk/packages/node/iii/src/types.ts、sdk/packages/python/iii/src/iii/iii.pyNode 的registerTriggerType会返回一个TriggerTypeRef带有.unregister()快捷方式内部委托给worker.unregisterTriggerType(...)Python 的TriggerTypeRef只暴露register_trigger和register_function拆除类型本身要走worker.unregister_trigger_type(...)Rust 只接收id字符串Node 与 Python 接收完整输入对象但实际只使用其中的id字段来定位被拆除的类型。向已绑定函数分发事件iii没有专门的触发API。当底层事件源送来内容一个 HTTP 请求、一次 cron 滴答、一次 webhook 命中时你的发布者 worker 在registerTrigger回调构建的bindings表中查找对应条目然后通过worker.trigger(...)调用每个匹配的函数。沿用上面的mini-http例子Node / TypeScript// Inside the workers HTTP listener, after matching methodpath to an // entry in the bindings map from the declare-trigger-type example: const binding bindings.get(matchedTriggerId); await worker.trigger({ function_id: binding.function_id, payload: { method, headers, body }, });Python# Inside the workers HTTP listener, after matching methodpath to an # entry in the bindings dict from the declare-trigger-type example: binding bindings[matched_trigger_id] worker.trigger({ function_id: binding.function_id, payload: {method: method, headers: headers, body: body}, })Rustuse iii_sdk::TriggerRequest; use serde_json::json; // Inside the workers HTTP listener, after matching methodpath to an // entry in the handlers bindings map from the declare-trigger-type example: let binding handler.bindings.lock().unwrap().get(matched_trigger_id).cloned(); if let Some(binding) binding { worker .trigger(TriggerRequest { function_id: binding.function_id.clone(), payload: json!({ method: method, headers: headers, body: body }), action: None, timeout_ms: None, }) .await?; }每一次分发事件时引擎都会评估消费者的config与可选的condition_function_id然后把匹配的调用路由到绑定函数并把结果返回给调用方。源码佐证引擎侧的注册与分发机制文档描述的引擎转发绑定/解绑回调、但不负责分发这一设计在引擎源码中有清晰的印证。TriggerRegistryengine/src/trigger.rs维护三类数据结构trigger_types按(namespace, type id)键控的提供者provider表triggers当前存活的绑定pending_triggers挂起意图——当触发器类型尚未注册、提供者 worker 断线、绑定未能送达或提供者异步拒绝激活时绑定会被停车park在这里它们被禁用不会触发直到触发器类型重新注册时被激活并移入triggers。这正是发布者未连接时注册也能成功的底层实现register_triggerengine/src/trigger.rs在找不到提供者时不会失败而是打印[PENDING]警告并把意图插入pending_triggers同时在类型注册的并发窗口内做一次 re-check 以关闭停车/排空竞争。register_trigger_typeengine/src/trigger.rs则负责先发布类型、再重放replay匹配的存活绑定、排空挂起意图并处理回迁re-homing——把曾回退到default命名空间的绑定迁移回新注册的自家提供者。unregister_workerengine/src/trigger.rs则展示了断线清理提供者离开后其绑定变成孤儿会重新解析提供者、尽力回退实在无处解析的绑定被[DISABLED]并停车等待类型回归时自动恢复。引擎对已知触发器类型的提供者做了内置映射KNOWN_TRIGGER_TYPE_PROVIDERS见 engine/src/trigger.rshttp、cron、subscribe、state、durable:subscriber、stream、log、trace、configuration等均在其中当挂起警告持久存在时日志会提示缺失的 worker 并给出iii trigger -n namespace compose::add workerworker的安装建议。小结iii 的触发器体系把事件路由的职责完全交给了发布者 worker引擎只负责记录绑定并转发回调如何查表、如何分发由你决定。实际开发中你通常只需用worker.registerTrigger消费既有类型当需要把自有的 HTTP、webhook、文件或数据库事件暴露给其他 worker 时则用registerTriggerTypeTriggerHandler声明类型、用worker.trigger分发事件并用trigger_request_format/call_request_format把契约写清楚。记住几个关键边界metadata 是消费者按绑定打的标签schema 是发布者按类型声明的契约且暂不参与运行时校验而触发器类型的生命周期断线自动清理、挂起重放、回迁由引擎的TriggerRegistry托管无需你手动兜底。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考