iii 触发器实战:注册、元数据、条件门控与反注册全解析
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
触发器(Trigger)是 iii 中连接事件源与函数的核心机制:函数不仅能被worker.trigger或iii trigger直接调用,还能在满足某个事件条件时被引擎自动唤起——例如一次http请求、一个cron定时任务、一次state状态变更,或任何 worker 发布的自定义事件。本文以 docs/using-iii/triggers.mdx 为主线,结合仓库内引擎与 SDK 源码,系统讲解触发器注册(绑定)、触发器元数据、触发器类型未就绪时的乐观注册、多触发器绑定同一函数、条件门控(gating)与反注册,并给出 Node/TypeScript、Python、Rust 三套可运行的完整代码。
阅读本文后,你将能够在自己的 iii 项目中熟练地把任意函数绑定到事件源、为绑定附加上下文元数据、用条件函数做流量闸门,并理解引擎在背后如何存储、延迟激活与回放这些绑定。
触发器是什么:把事件绑定到函数
在 iii 中,"注册一个触发器"(register a trigger)与"注册一个触发器类型"(register a trigger type)是两件不同的事,本文聚焦前者:
- 注册触发器(绑定):消费者(consumer)把某个触发器类型(如
http)绑定到自己的某个函数上,声明"当这类事件发生时要调用哪个函数"; - 注册触发器类型:发布者(publisher)在 worker 里声明一种新的事件类型,让其他 worker 的函数可以绑定到它,详见 Creating Workers / Triggers。
绑定通过function_id完成。一个触发器声明三件事:
| 字段 | 含义 |
|---|---|
type | 触发器类型标识,如http、cron、state、durable:subscriber |
config | 由每个触发器类型自行定义的结构化配置,例如http的api_path与http_method |
function_id | 触发器触发时要调用的目标函数 |
从引擎源码看,绑定最终会落成一条Trigger记录:它包含id、trigger_type、function_id、config、可选的metadata,以及一组用于命名空间解析的字段(namespace、trigger_namespace、home_namespace、provider_namespace),见 engine/src/trigger.rs。其中config的类型和字段由各发布 worker 定义,引擎为内置类型提供了类型化定义(下文有对照表)。
注册一个触发器:三种 SDK 的完整示例
下面用同一份"把math::add绑定到 HTTP POST 端点"的需求,展示三个 SDK 的注册写法。
Node / TypeScript(见 sdk/packages/node/iii/src/iii.ts 的registerTrigger实现):
import { registerWorker } from "iii-sdk"; const url = process.env.III_URL; if (!url) throw new Error("III_URL must be set"); const worker = registerWorker(url); worker.registerTrigger({ type: "http", function_id: "math::add", config: { api_path: "/math/add", http_method: "POST" }, });Python:
import os from iii import register_worker, InitOptions worker = register_worker( os.environ.get("III_URL"), InitOptions(worker_name="my-worker"), ) worker.register_trigger({ "type": "http", "function_id": "math::add", "config": {"api_path": "/math/add", "http_method": "POST"}, })Rust:
use iii_sdk::{InitOptions, RegisterTriggerInput, register_worker}; use serde_json::json; let url = std::env::var("III_URL").expect("III_URL must be set"); let worker = register_worker(&url, InitOptions::default()); worker.register_trigger(RegisterTriggerInput { trigger_type: "http".into(), function_id: "math::add".into(), config: json!({ "api_path": "/math/add", "http_method": "POST" }), metadata: None, })?;几个值得注意的细节:
- Node SDK 的
registerTrigger在本地用crypto.randomUUID()生成触发器id,然后通过 WebSocket 发送RegisterTrigger消息给引擎,同时把完整触发器存入本地this.triggers表,返回一个带unregister()方法的句柄(详见下文"反注册"一节)。 - 命名空间上,Node SDK 会做一次"未显式声明则默认本 worker 命名空间"的处理:
callNamespace(trigger.namespace, this.namespace, ...)——触发器点名一个函数,而函数注册在 worker 自己的命名空间里,如果触发器默认落在引擎的default命名空间,就会"触发却解析不到函数"(见 sdk/packages/node/iii/src/iii.ts)。 - 内置触发器的
config形状由引擎统一生成 JSON Schema。例如http的注册配置包含api_path(如/users/:id)、http_method(默认GET)以及可选的condition_function_id,见 engine/src/trigger_formats.rs。
各类型完整的注册配置与触发负载(call request)形状,可在引擎的内置类型表中查到:http、cron、subscribe、state、durable:subscriber、stream、stream:join、stream:leave、log、trace、configuration均有类型化定义,见 engine/src/trigger_formats.rs。
触发器元数据:给绑定附加任意上下文
注册触发器时可以带上可选的metadata字段(上面示例里的null/None)。它是一段随触发器一起存储的任意 JSON,在触发器触发时,会作为一个独立参数与 payload 一起交给目标函数——而不是混在 payload 里。
它的典型用途是提供触发上下文:当一个函数被多个触发器共享时,处理函数可以用metadata反推"这次是哪个注册项触发的、带着什么上下文"。
// Node / TypeScript:带上团队与环境标签 worker.registerTrigger({ type: "http", function_id: "math::add", config: { api_path: "/math/add", http_method: "POST" }, metadata: { team: "platform", env: "staging" }, });# Python worker.register_trigger({ "type": "http", "function_id": "math::add", "config": {"api_path": "/math/add", "http_method": "POST"}, "metadata": {"team": "platform", "env": "staging"}, })// Rust 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" })), })?;底层事实链非常清晰:
- 引擎的
Trigger结构体把metadata序列化进绑定记录(#[serde(skip_serializing_if = "Option::is_none")]),见 engine/src/trigger.rs; - 触发分发时,
Engine::fire_triggers取出该触发器存储的metadata,以独立参数调用目标函数:call_with_metadata_ns(&namespace, &function_id, data, metadata),见 engine/src/engine/mod.rs; - SDK 协议层同样支持:Python 的测试 sdk/packages/python/iii/tests/test_trigger_metadata.py 验证了
RegisterTriggerInput与RegisterTriggerMessage均能携带并原样传递metadata。
注意两点:其一,metadata既可以通过registerTrigger提供,也可以在直接trigger()调用时提供;其二,触发器类型本身没有 metadata 字段,metadata 是挂在每次绑定上的,与触发器类型声明的 schema(trigger_request_format/call_request_format)不是一回事——前者由消费者在绑定时设置,是发布者做台账与发现用的自由标签;后者由发布者声明,用于描述消费者该传什么、会收到什么。关于 handler 侧如何读取每次调用的元数据,可参考 Creating Workers / Functions 中"Receive per-invocation metadata"一节。
触发器类型尚未就绪时:乐观注册与延迟激活
iii 的触发器注册是与顺序无关的。如果注册时该触发器类型在项目中还没有激活(例如发布http事件的 worker 尚未连接),引擎不会报错,而是把注册"乐观地"存下来,等该触发器类型一上线就自动激活。这一机制在引擎中有一套完整的落盘与恢复逻辑:
- 引擎的
TriggerRegistry维护两张映射:triggers(活跃绑定)与pending_triggers(等待激活的绑定意图),见 engine/src/trigger.rs; register_trigger在找不到可用 provider 时,会打印[PENDING]警告并把绑定插入pending_triggers,返回RegisterTriggerOutcome::Deferred;只有在 provider 存在且注册成功时才返回Registered,见 engine/src/trigger.rs;- 当某个触发器类型(重新)注册时,
register_trigger_type会回放(replay)已存在的绑定,并逐个把pending_triggers里等待该类型的意图取出、激活、移入triggers,见 engine/src/trigger.rs; - 发布 worker 重启后,之前存储的绑定会被重新建立,因此"消费者先起来、发布者后起来"的顺序完全可以正常工作。
下面是一段可完整运行的演示:先在httpworker 未启动的情况下注册绑定。
// Node / TypeScript —— http worker not started worker.registerTrigger({ type: "http", function_id: "math::add", config: { api_path: "/math/add", http_method: "POST" }, });# Python —— http worker not started worker.register_trigger({ "type": "http", "function_id": "math::add", "config": {"api_path": "/math/add", "http_method": "POST"}, })// Rust —— http worker not started worker.register_trigger(RegisterTriggerInput { trigger_type: "http".into(), function_id: "math::add".into(), config: json!({ "api_path": "/math/add", "http_method": "POST" }), metadata: None, })?;随后再启动httpworker。引擎会自动激活已存储的绑定,无需重新注册,端点即可服务请求:
# http worker started after registration iii worker add http # the stored binding is now live on the http worker's default port (3111) curl -X POST http://localhost:3111/math/add \ -H 'Content-Type: application/json' \ -d '{"a": 2, "b": 3}' # 200 OK: the request reaches math::add through the activated trigger对已知触发器类型(如http、state、durable:subscriber、stream),引擎的[PENDING]提示还会附带提供该类型的 worker 的安装指引。引擎内置了一张"触发器类型 → 提供 worker"的映射表,见 engine/src/trigger.rs:
| 触发器类型 | 提供 worker |
|---|---|
http | http |
cron | cron |
subscribe | pubsub |
state | state |
durable:subscriber | queue |
stream/stream:join/stream:leave | iii-stream |
log/trace | iii-observability |
configuration | configuration |
pending_trigger_warning会基于这张表生成可执行建议(例如缺失httpworker 时提示安装命令),见 engine/src/trigger.rs。
什么时候注册会失败?当活跃的 provider 拒绝了给定的config(例如触发器配置非法)时,注册才会失败:引擎会把TriggerRegistrationResult(带errorbody)发回发起注册的 worker 并记录日志。对应地,引擎在register_trigger里对 provider 的拒绝做了区分:只有"provider 已处理并拒绝"才是确定性的失败;若只是连接通道不可达(worker 正在断开),则视为投递失败,把绑定暂存为 pending,等类型(重新)注册时再重试,见 engine/src/trigger.rs。
一个函数绑定多个触发器
同一个function_id可以绑定任意数量的触发器,且可以横跨不同类型的触发器。绑定第二个触发器只需复用相同的function_id,换一个类型或配置即可;函数本身无需改动——无论是来自 HTTP 请求、cron 定时还是队列消息,函数都以同样的方式被调用。
下面把reports::generate同时绑到 HTTP POST 与每周一次的 cron 上:
// Node / TypeScript —— Same handler runs for an HTTP POST and a weekly cron tick. worker.registerTrigger({ type: "http", function_id: "reports::generate", config: { api_path: "/reports/generate", http_method: "POST" }, }); worker.registerTrigger({ type: "cron", function_id: "reports::generate", config: { expression: "0 0 9 * * 1" }, // Every Monday at 09:00 });# Python worker.register_trigger({ "type": "http", "function_id": "reports::generate", "config": {"api_path": "/reports/generate", "http_method": "POST"}, }) worker.register_trigger({ "type": "cron", "function_id": "reports::generate", "config": {"expression": "0 0 9 * * 1"}, # Every Monday at 09:00 })// Rust worker.register_trigger(RegisterTriggerInput { trigger_type: "http".into(), function_id: "reports::generate".into(), config: json!({ "api_path": "/reports/generate", "http_method": "POST" }), metadata: None, })?; worker.register_trigger(RegisterTriggerInput { trigger_type: "cron".into(), function_id: "reports::generate".into(), config: json!({ "expression": "0 0 9 * * 1" }), // Every Monday at 09:00 metadata: None, })?;cron类型的expression采用 6 段格式(秒 分 时 日 月 星期),引擎侧定义见 engine/src/trigger_formats.rs。触发时,cron负载会携带trigger、job_id、scheduled_time、actual_time等字段(engine/src/trigger_formats.rs)。
用条件函数门控触发器
触发器可以在config中携带可选的condition_function_id。当触发器触发时,引擎会先用原本要传给 handler 的同一个 payload 调用条件函数;只有当条件函数返回真值时,目标function_id才会执行。条件函数就是一个普通的已注册函数。
引擎的实现位于 engine/src/condition.rs:条件函数在与触发器目标函数相同的命名空间内解析执行,返回语义为——Ok(true)放行、Ok(false)跳过、调用失败则整体报错;若条件函数返回None或非布尔值,check_condition也会放行(result.as_bool() != Some(false))。这一行为在 engine/src/condition.rs 的单测中逐一验证,并被内置的 queue、stream、configuration 等 worker 在分发路径上调用(例如 engine/src/workers/queue/adapters/builtin/adapter.rs)。
下面实现一个"仅黄金会员订单走加急通道"的门控:
// Node / TypeScript worker.registerFunction( "orders::is-priority", async (payload: { customer_tier: string }) => payload.customer_tier === "gold", ); worker.registerTrigger({ type: "http", function_id: "orders::expedite", config: { api_path: "/orders/expedite", http_method: "POST", condition_function_id: "orders::is-priority", }, });# Python def is_priority(payload: dict) -> bool: return payload.get("customer_tier") == "gold" worker.register_function("orders::is-priority", is_priority) worker.register_trigger({ "type": "http", "function_id": "orders::expedite", "config": { "api_path": "/orders/expedite", "http_method": "POST", "condition_function_id": "orders::is-priority", }, })// Rust use iii_sdk::{RegisterFunction, RegisterTriggerInput}; use schemars::JsonSchema; use serde::Deserialize; use serde_json::json; #[derive(Deserialize, JsonSchema)] struct Payload { customer_tier: String } worker.register_function(RegisterFunction::new( "orders::is-priority", |input: Payload| -> Result<bool, String> { Ok(input.customer_tier == "gold") }, )); worker.register_trigger(RegisterTriggerInput { trigger_type: "http".into(), function_id: "orders::expedite".into(), config: json!({ "api_path": "/orders/expedite", "http_method": "POST", "condition_function_id": "orders::is-priority", }), metadata: None, })?;从源码结构看,condition_function_id被内置在多个触发器类型的配置结构里(http、cron、state、stream、queue、configuration等的配置 struct 都声明了该可选字段),因此门控能力并不只属于 HTTP,而是事件分发路径上的通用机制,见 engine/src/trigger_formats.rs。
反注册触发器
registerTrigger(及各语言等价调用)会返回一个带unregister()方法的句柄。调用它即可在运行时移除触发器;而当 worker 断开连接时,它所注册的所有触发器会被自动清理。
// Node / TypeScript const trigger = worker.registerTrigger({ type: "http", function_id: "math::add", config: { api_path: "/math/add", http_method: "POST" }, }); trigger.unregister();# Python trigger = worker.register_trigger({ "type": "http", "function_id": "math::add", "config": {"api_path": "/math/add", "http_method": "POST"}, }) trigger.unregister()// Rust let trigger = worker.register_trigger(RegisterTriggerInput { trigger_type: "http".into(), function_id: "math::add".into(), config: json!({ "api_path": "/math/add", "http_method": "POST" }), metadata: None, })?; trigger.unregister();引擎侧的反注册是幂等的:unregister_trigger对不存在的 id 返回Ok(false)而非报错,重复反注册等价于空操作;如果该绑定仍处于 pending 状态,反注册等价于丢弃这个等待中的意图。引擎会先通知 provider 的 registrator 解除绑定,再在triggers与pending_triggers中清除该 id(见 engine/src/trigger.rs)。Node SDK 的反注册则通过发送UnregisterTrigger消息并删除本地句柄实现(见 sdk/packages/node/iii/src/iii.ts)。worker 断线时的批量清理在TriggerRegistry::unregister_worker中完成:它会移除该 worker 拥有的触发器类型与绑定,并把失去 provider 的绑定重新解析或暂存为 pending,确保 provider 重启不会静默丢掉其他人的绑定(engine/src/trigger.rs)。
总结:触发器生命周期与源码地图
一个触发器从注册到销毁的完整生命周期:
- 注册:SDK 发送
RegisterTrigger消息,引擎把Trigger(含type、config、function_id、metadata)存入注册表; - 激活:provider 在线则立即激活并通知 registrator;provider 未上线则进入
pending_triggers,待类型(重新)注册时自动回放激活; - 触发:事件发生时引擎按类型收集绑定,把 payload 与
metadata作为独立参数调用目标函数;若配置了condition_function_id,先求值条件函数再决定是否放行; - 反注册:显式
unregister()或 worker 断线自动清理,均幂等。
想深入验证或二次开发,可在仓库中按如下地图继续阅读:
- 引擎注册表与乐观注册/回放/再归属逻辑:engine/src/trigger.rs
- 内置触发器类型的配置与负载 schema:engine/src/trigger_formats.rs
- 条件函数求值语义:engine/src/condition.rs
- 触发分发(含 metadata 独立传参):engine/src/engine/mod.rs
- Node SDK 的
registerTrigger/registerTriggerType与反注册句柄:sdk/packages/node/iii/src/iii.ts - Node SDK 的
TriggerConfig/TriggerHandler类型定义:sdk/packages/node/iii/src/triggers.ts - 编写自定义触发器类型(发布者视角):docs/creating-workers/triggers.mdx
- 直接调用函数(
worker.trigger/iii trigger)与engine::triggers::list发现能力:docs/using-iii/functions.mdx
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考