iii Workers 接入指南:从脚手架到优雅下线的完整生命周期
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
本指南以 iii 项目docs/0-17-0/creating-workers/workers.mdx文档为主体,系统讲解如何将一个 Worker 部署并接入 iii Engine:从iii worker init脚手架生成、WebSocket 连接、生命周期状态机、注册表巡检、断线容错与发现事件订阅,到iii.worker.yaml清单配置与shutdown优雅退出。读完本文,你将掌握用 Node/TypeScript、Python、Rust 三种 SDK 编写并运维一个可被 Engine 路由调用的 Worker 的完整实战方案。
Workers 如何扩展 iii
Worker 是 iii 系统中能力的载体。每一个 Worker 都向系统贡献两类可路由资源,Engine 可以据此将调用与事件分发到正确的目标:
- Functions:通过
function_id从系统任意位置调用(见 docs/0-17-0/using-iii/functions.mdx); - Triggers:Worker 主动通告(advertise)的触发器类型,其他 Worker 可以将自己的 Function 绑定到这些触发器上(见 docs/0-17-0/using-iii/triggers.mdx)。
关于 Worker 与 iii 交互所需的完整 SDK 表面(API),可查阅对应语言的 SDK 参考:Node、Python、Rust、Browser。
从引擎内部看,"worker" 这一抽象贯穿了注册、发现、调用全链路:engine::*::list系列内建 Function 与engine::workers-available/engine::functions-available发现事件正是由 Engine 内建 worker 暴露的(参见 engine/src/workers/engine_fn/README.md),后续章节会逐一展开。
脚手架:iii worker init新建 Worker
iii worker init从零创建一个独立的 Worker 工程。该命令会写入一个语言专属的项目目录,其中包含已安装的 iii SDK、一份iii.worker.yaml清单,以及可供替换的示例 Function 与 Trigger 注册代码。
# 交互式:提示选择语言 iii worker init my-worker # 全脚本化:传入 --language 跳过交互提示 iii worker init my-worker --language typescript支持的脚手架语言:
| 语言 | 别名 |
|---|---|
typescript | ts |
javascript | js |
python | py |
rust | rs |
关于该命令的几点行为细节,务必留意:
- 位置参数
NAME是目标目录名;可通过--directory覆盖。 - 在已经包含 iii Worker 的目录(即存在
.iii/worker.ini)上重复执行iii worker init,不会对既有 Worker 做任何修改——.iii/worker.ini是 Worker 已初始化的标记文件。 - 默认情况下,对非空目录执行 init 会失败;如需向其他非空目录脚手架,使用
--allow-non-empty。
若想从注册表安装一个已存在的 Worker,而不是新脚手架一个,请使用
iii worker add。注册表的使用方式见 docs/0-17-0/using-iii/workers.mdx 的 Finding Workers 小节。
从源码结构看,CLI 侧的 Worker 相关能力集中在iii-worker二进制与worker子命令路径上(见 engine/src/cli/registry.rs),脚手架产出的目录结构、SDK 安装与示例注册代码均可在仓库 fixtures 与集成测试(如 engine/tests/builtin_functions_e2e.rs)中看到对应用法。
连接 Engine:WebSocket 与III_URL
Worker 通过WebSocket连接 Engine。约定做法是通过环境变量III_URL设置 Engine 地址,也可以显式传给register_worker。连接字符串是 Worker 与它加入的 iii 实例之间的唯一耦合点,因此Worker 进程可以部署在网络可达的任何位置——这为跨主机、跨集群部署提供了自由。
三种语言 SDK 的接入方式:
// Node / TypeScript 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, { workerName: "my-worker", });# Python import os from iii import register_worker, InitOptions worker = register_worker( os.environ.get("III_URL"), InitOptions(worker_name="my-worker"), )// Rust use iii_sdk::{InitOptions, WorkerMetadata, register_worker}; let url = std::env::var("III_URL").expect("III_URL must be set"); let worker = register_worker( &url, InitOptions { metadata: Some(WorkerMetadata { name: "my-worker".into(), ..Default::default() }), ..Default::default() }, );从 SDK 的工程实践看,workerName/worker_name/WorkerMetadata.name会作为注册标识出现在 Engine 的注册表中,是发现与巡检(见下文engine::workers::list)时区分 Worker 的关键字段。
Worker 生命周期
状态机
Worker 连接后会在一个很小的状态集合间流转:
connecting → connected → available / busy → disconnectedconnecting:WebSocket 握手阶段;connected:Worker 已加入 Engine 的注册表;available/busy:描述 Worker 当前是否正在处理调用(busy 表示正在处理中);disconnected:WebSocket 关闭后的终态。
Engine 会跟踪这些状态转换,并通过其发现函数(discovery functions)把它们呈现给其他 Worker 和工具,使系统其余部分能够对这些变化做出反应。
巡检实时注册表
要查看当前连接到 Engine 的内容,调用engine::*::list系列内建 Function 获取注册表当前状态。每个 Function 返回一个列表:
| Function | 返回值 |
|---|---|
engine::workers::list | 每个已连接 Worker 及其指标。 |
engine::functions::list | 每个已注册 Function,可通过include_internal过滤。 |
engine::triggers::list | 每个已注册 Trigger,可通过include_internal过滤。 |
engine::trigger-types::list | 每个已通告的 Trigger 类型及其配置与调用 schema。 |
三语言示例(可用{ "worker_id": "<uuid>" }精确查询单个 Worker):
// Node / TypeScript const { workers } = await worker.trigger({ function_id: "engine::workers::list", payload: {}, }); const { functions } = await worker.trigger({ function_id: "engine::functions::list", payload: { include_internal: false }, }); const { triggers } = await worker.trigger({ function_id: "engine::triggers::list", payload: { include_internal: false }, }); const { trigger_types } = await worker.trigger({ function_id: "engine::trigger-types::list", payload: { include_internal: false }, });# Python workers = worker.trigger({ "function_id": "engine::workers::list", "payload": {}, })["workers"] functions = worker.trigger({ "function_id": "engine::functions::list", "payload": {"include_internal": False}, })["functions"] triggers = worker.trigger({ "function_id": "engine::triggers::list", "payload": {"include_internal": False}, })["triggers"] trigger_types = worker.trigger({ "function_id": "engine::trigger-types::list", "payload": {"include_internal": False}, })["trigger_types"]// Rust use iii_sdk::TriggerRequest; use serde_json::json; let workers = worker .trigger(TriggerRequest { function_id: "engine::workers::list".into(), payload: json!({}), action: None, timeout_ms: None, }) .await?; let functions = worker .trigger(TriggerRequest { function_id: "engine::functions::list".into(), payload: json!({ "include_internal": false }), action: None, timeout_ms: None, }) .await?; let triggers = worker .trigger(TriggerRequest { function_id: "engine::triggers::list".into(), payload: json!({ "include_internal": false }), action: None, timeout_ms: None, }) .await?; let trigger_types = worker .trigger(TriggerRequest { function_id: "engine::trigger-types::list".into(), payload: json!({ "include_internal": false }), action: None, timeout_ms: None, }) .await?;include_internal参数用于决定是否把 Engine 内部 Worker 的注册项纳入结果;内部 Worker(如 telemetry worker)并不面向用户配置,也没有iii.worker.yaml(参见 engine/src/workers/telemetry/README.md),巡检时默认过滤掉它们更符合业务视角。
处理 Worker 断线
当 Worker 的 WebSocket 关闭时,Engine 会自动为其清理:它的 Functions 和 Triggers 会从实时注册表中移除,任何针对这些 Function 的在途调用(in-flight invocations)都会被取消。
在途请求:捕获invocation_stopped
在途请求会收到invocation_stopped错误。请捕获这类错误并按"取消"处理。在拥有该 Function 的 Worker 重连之前,重试都会失败。
// Node / TypeScript import { IIIInvocationError } from "iii-sdk"; try { const result = await worker.trigger({ function_id: "math::add", payload: { a: 1, b: 2 }, }); } catch (err) { if (err instanceof IIIInvocationError && err.code === "invocation_stopped") { // Worker 在调用中途断开。订阅 `engine::functions-available` // (见下文 "Subscribe to changes")以获知何时可以重试。 return; } throw err; }# Python from iii import IIIInvocationError try: result = worker.trigger({ "function_id": "math::add", "payload": {"a": 1, "b": 2}, }) except IIIInvocationError as err: if err.code == "invocation_stopped": # Worker 在调用中途断开。订阅 `engine::functions-available` # (见下文 "Subscribe to changes")以获知何时可以重试。 return raise// Rust use iii_sdk::{IIIError, TriggerRequest}; use serde_json::json; let result = worker .trigger(TriggerRequest { function_id: "math::add".into(), payload: json!({ "a": 1, "b": 2 }), action: None, timeout_ms: None, }) .await; match result { Err(IIIError::Remote { code, .. }) if code == "invocation_stopped" => { // Worker 在调用中途断开。订阅 `engine::functions-available` // (见下文 "Subscribe to changes")以获知何时可以重试。 } Err(e) => return Err(e.into()), Ok(value) => { /* use value */ } }底层行为与引擎侧的实现一致:Engine 在 Worker 断开时会立即驱逐该 Worker 的 Functions,并将其在途调用以invocation_stopped解析(见 engine/src/workers/engine_fn/skills/SKILL.md)。
订阅变更:发现事件
你可以把 Trigger 注册到 Engine 的发现事件上,实时响应拓扑变化。这在 Worker 恢复在线后继续未完成工作时尤其有用。
| Trigger | 触发时机 |
|---|---|
engine::workers-available | 有 Worker 连接或断开。 |
engine::functions-available | 有 Function 注册或注销。 |
// Node / TypeScript worker.registerFunction( "discovery::on-workers", async (data: { event: string; worker_id: string }) => { if (data.event === "worker_connected") { // 有 Worker 刚加入注册表,它的 Functions 现在可调用了。 } }, ); worker.registerTrigger({ type: "engine::workers-available", function_id: "discovery::on-workers", config: {}, }); worker.registerFunction( "discovery::on-functions", async (data: { event: string; functions: { function_id: string }[] }) => { // `functions` 是变更后的完整快照。 const ids = data.functions.map((f) => f.function_id); }, ); worker.registerTrigger({ type: "engine::functions-available", function_id: "discovery::on-functions", config: {}, });# Python async def on_workers(data: dict) -> None: if data["event"] == "worker_connected": # 有 Worker 刚加入注册表,它的 Functions 现在可调用了。 pass worker.register_function("discovery::on-workers", on_workers) worker.register_trigger({ "type": "engine::workers-available", "function_id": "discovery::on-workers", "config": {}, }) async def on_functions(data: dict) -> None: # `functions` 是变更后的完整快照。 ids = [f["function_id"] for f in data.get("functions", [])] worker.register_function("discovery::on-functions", on_functions) worker.register_trigger({ "type": "engine::functions-available", "function_id": "discovery::on-functions", "config": {}, })// Rust use iii_sdk::{RegisterFunction, RegisterTriggerInput}; use schemars::JsonSchema; use serde::Deserialize; use serde_json::{Value, json}; #[derive(Deserialize, JsonSchema)] struct WorkersAvailable { event: String, worker_id: String } #[derive(Deserialize, JsonSchema)] struct FunctionsAvailable { event: String, functions: Vec<Value> } worker.register_function(RegisterFunction::new_async( "discovery::on-workers", |input: WorkersAvailable| async move { if input.event == "worker_connected" { // 有 Worker 刚加入注册表,它的 Functions 现在可调用了。 } Ok::<_, String>(()) }, )); worker.register_trigger(RegisterTriggerInput { trigger_type: "engine::workers-available".into(), function_id: "discovery::on-workers".into(), config: json!({}), metadata: None, })?; worker.register_function(RegisterFunction::new_async( "discovery::on-functions", |input: FunctionsAvailable| async move { // `functions` 是变更后的完整快照。 let _count = input.functions.len(); Ok::<_, String>(()) }, )); worker.register_trigger(RegisterTriggerInput { trigger_type: "engine::functions-available".into(), function_id: "discovery::on-functions".into(), config: json!({}), metadata: None, })?;需要说明的是,SDK 具备自动重连能力(带退避),并会原样重放注册项(replays registrations verbatim),因此不要手动重复注册。调用方在断线窗口内看到的是invocation_stopped——请把它当作"取消"语义,而不是瞬时故障(见 engine/src/workers/engine_fn/skills/SKILL.md)。发现事件 + 自动重连的组合,正是构建自愈型 Worker 拓扑的关键。
Worker 清单:iii.worker.yaml
iii.worker.yaml是位于 Worker 根目录的清单文件,它告诉 iii 如何安装依赖、运行 Worker 以及透传配置。它同时适用于两类场景:
- iii worker CLI 命令(如
start、stop、restart,见 docs/0-17-0/using-iii/workers.mdx 的 Starting and Stopping Workers 小节); - 当 Worker 被声明在 iii 的
config.yaml中时,由 iii 自动启动的 Worker。
name: math-worker runtime: kind: python package_manager: pip entry: math_worker.py scripts: install: "pip install -r requirements.txt" start: "python math_worker.py"字段要点:
name:Worker 名称,同时被用作注册表 / 清单中的标识(参见 engine/src/workers/engine_fn/README.md 中iii.worker.yamlname:字段的说明);runtime.kind:运行时类型(如python),runtime.package_manager指定依赖管理器(如pip),runtime.entry指定入口文件;scripts.install:安装依赖的命令;scripts.start:启动 Worker 的命令。
清单只是关于启动Worker 的元数据。一旦 Worker 运行起来,iii 对它们一视同仁:由config.yaml启动的 Worker、iii worker start启动的 Worker,以及手动运行、直接使用 iii SDK 的进程,与 Engine 的交互行为完全一致。此外,iii.worker.yaml中也可以承载 Worker 自身的配置块(例如 observability 相关的配置项,参见 engine/src/workers/observability/config.rs 中对iii.worker.yamlconfig block 的处理逻辑)。
如果 Worker 无法正常启动,请检查其清单,并使用
iii worker logs查看 Worker 日志定位问题。
关闭 Worker:优雅下线与一次性 Worker
调用 SDK 的shutdown可以干净地关闭 WebSocket。此时 Engine 会:
- 将 Worker 的 Functions 和 Triggers 从注册表中移除;
- 触发
engine::workers-available,事件为worker_disconnected; - 以
invocation_stopped取消针对这些 Function 的在途调用。
即使不调用shutdown,进程突然退出也会在 Engine 发现 socket 断开后到达相同的状态;但优雅下线让这一过程确定且更快——SDK 侧的shutdown会先冲刷待发送流量再关闭 WebSocket(见 engine/src/workers/engine_fn/skills/SKILL.md)。
三语言示例:
// Node / TypeScript process.on("SIGTERM", async () => { await worker.shutdown(); process.exit(0); });# Python import signal def _on_term(*_): worker.shutdown() raise SystemExit(0) signal.signal(signal.SIGTERM, _on_term)// Rust // Rust 线程本身不会维持进程存活;在 `main` 返回前 await 此调用, // 以便连接线程干净退出。 worker.shutdown_async().await;
shutdown对一次性(One-shot / ephemeral)Worker尤其有用。Kubernetes Job、Serverless 容器或定时脚本可以像任何其他 Worker 一样连接,完成工作后调用shutdown()(Rust 中为shutdown_async().await)干净退出,不留残余连接。
小结
Worker 是 iii 系统的能力单元,其生命周期可以概括为一条完整链路:iii worker init脚手架 → 通过III_URL建立 WebSocket 连接 → 进入connecting → connected → available / busy → disconnected状态机 → 用engine::*::list巡检注册表、用发现事件响应拓扑变化、用invocation_stopped语义处理断线 → 由iii.worker.yaml定义启动方式 → 最终以shutdown优雅下线。这套机制让 Worker 既能常驻服务,也能作为一次性任务短暂接入,而 Engine 侧的自动清理与发现事件保证了整个系统在 Worker 频繁加入、离开时依然稳定可观测。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考