iii Worker 完全指南:从脚手架生成到生命周期管理(附源码实现剖析)
iii Worker 完全指南从脚手架生成到生命周期管理附源码实现剖析【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii本文基于 iii 官方文档0-16.0 版本中「Creating Workers / Workers」页面整理与扩充系统讲解如何用iii worker init脚手架生成 Worker、通过 WebSocket 将其接入 Engine、利用引擎发现discovery接口检查注册表、处理断连与优雅关闭等完整链路。读完本文你将掌握从源码级参数语义到 SDK 调用模式的全部实操知识并能在多语言Node/TypeScript、Python、Rust环境下构建并稳定运维 iii Worker。Worker 为 iii 系统带来什么Worker 是向 iii 系统扩展能力的基本单元。每个 Worker 向 Engine 贡献一组函数Functions与一组触发器TriggersEngine 负责把调用路由到对应的 Worker。接入后一个 Worker 对外暴露两类资源Functions系统中任何位置都可以通过function_id调用Triggers它声明的触发类型其他 Worker 可以把自身函数绑定bind上去。需要完整的 SDK 能力面各语言 SDK 的完整 API时可参考仓库文档 sdk-reference 目录下的 Node、Python、Rust、Browser 各语言参考以及 Using iii / Functions、Using iii / Triggers。从源码结构看Engine 侧对 Worker 的管理集中在 engine/src/workers/registry.rs注册表、engine/src/workers/traits.rsWorker 抽象接口与 engine/src/workers/engine_fn/mod.rs引擎内建函数与触发器三个模块这正是下文检查注册表与订阅变更两节所有接口背后的实现位置。用iii worker init脚手架生成新 Workeriii worker init从零创建一个独立 Worker。命令会写出一个语言专属的项目目录其中已经装好 iii SDK、包含一个iii.worker.yaml清单以及可直接替换的示例函数与触发器注册代码。# 交互模式提示选择语言 iii worker init my-worker # 完全脚本化传 --language 跳过交互 iii worker init my-worker --language typescript支持的语言及短别名经 crates/iii-worker/src/cli/init.rs 中parse_language_arg验证语言长名短别名TypeScripttypescripttsJavaScriptjavascriptjsPythonpythonpyRustrustrs关键参数与行为均已与 CLI 源码逐项核对位置参数NAME是目标目录--directory优先级更高见 InitArgs 定义。两者都不传时目录默认为当前目录。对已存在 iii Worker 的目录即含.iii/worker.ini重复执行iii worker init不会产生任何变更——这是一次幂等的重新初始化源码注释中明确标注为 idempotent re-init。默认情况下指向非空目录会直接失败需要--allow-non-empty才允许在任意非空目录中脚手架。各语言默认入口文件不同default_entryTS/JS 为./src/index.ts/./src/index.jsPython 为./src/main.pyRust 为./src/main.rsinit 成功提示会引导你编辑这些文件。脚手架底层复用scaffolder_core的worker-bare模板模板中声明的renames规则会在拷贝时把iii.worker.lang.yaml、package.lang.json统一重命名为iii.worker.yaml、package.json。提示如果是从注册表安装已有的 Worker 而不是新建应使用iii worker add注册表能力面见 Using iii / Workers。连接 Engine唯一耦合点是连接串Worker 通过WebSocket连接 Engine。惯例是用III_URL环境变量提供 Engine 地址也可以把它显式传给register_worker。连接串是 Worker 与所加入 iii 实例之间唯一的耦合因此 Worker 进程可以部署在网络上任何可达的位置。Node / TypeScriptimport { 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, });Pythonimport os from iii import register_worker, InitOptions worker register_worker( os.environ.get(III_URL), InitOptions(worker_namemy-worker), )Rustuse 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() }, );Worker 生命周期状态机Worker 连接后会经历一组小型状态机connecting → connected → available / busy → disconnected。connectingWebSocket 握手阶段connectedWorker 已加入 Engine 的注册表available/busy描述 Worker 当前是否正在处理调用disconnectedWebSocket 关闭时的终态。Engine 会跟踪这些状态迁移并通过其发现函数把变更暴露给其他 Worker 与工具链让整个系统可以做出反应。检查实时注册表要查看当前有哪些实体连在 Engine 上调用engine::*::list系列函数即可每个都返回列表Function返回内容engine::workers::list所有已连接 Worker 及其指标engine::functions::list所有已注册 Function可按include_internal过滤engine::triggers::list所有已注册 Trigger可按include_internal过滤engine::trigger-types::list所有已声明的 Trigger 类型含配置与调用 schema这三个内置函数在引擎中确实是内建能力engine/src/engine/mod.rs 中有is_engine_owned_builtin(..., engine::functions::list)的断言与内建函数注册逻辑而worker_id的公开性即engine::workers::list可以按 UUID 精确查找单个 Worker也在 engine/src/worker_connections/mod.rs 的注释中得到确认。Node / TypeScript 示例其余语言同理payload 结构一致// engine::workers::list传 { worker_id: uuid } 可查询单个 worker const { workers } await worker.trigger({ function_id: engine::workers::list, payload: {}, }); // engine::functions::list const { functions } await worker.trigger({ function_id: engine::functions::list, payload: { include_internal: false }, }); // engine::triggers::list const { triggers } await worker.trigger({ function_id: engine::triggers::list, payload: { include_internal: false }, }); // engine::trigger-types::list 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 侧使用TriggerRequest结构体use iii_sdk::TriggerRequest; use serde_json::json; // engine::workers::list传 json!({ worker_id: uuid }) 可查询单个 worker let workers worker .trigger(TriggerRequest { function_id: engine::workers::list.into(), payload: json!({}), action: None, timeout_ms: None, }) .await?; // engine::functions::listtriggers / trigger-types 同构function_id 与返回值字段相应替换 let functions worker .trigger(TriggerRequest { function_id: engine::functions::list.into(), payload: json!({ include_internal: false }), action: None, timeout_ms: None, }) .await?;处理 Worker 断连当某个 Worker 的 WebSocket 关闭时Engine 会自动清理它的 Functions 与 Triggers 离开实时注册表指向这些 Function 的在途调用被取消。在途请求在途请求会收到invocation_stopped错误。应捕获该错误并当作取消处理在拥有该函数的 Worker 重新连接之前重试都会失败。这一错误码在引擎源码中确有定义与测试engine/src/invocation/mod.rs 构造了code: invocation_stopped的错误同文件 L433 附近有assert_eq!(error.code, invocation_stopped)的单元测试断言。引擎内建技能文档 engine/src/workers/engine_fn/skills/SKILL.md 也明确写道SDK 会以退避策略自动重连并原样重放注册——不要手动重复注册调用方在断连窗口内看到的invocation_stopped应视为取消而非瞬态故障。捕获示例Node / TypeScriptimport { 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 // 见下文订阅变更来判断何时重试。 return; } throw err; }Pythonfrom 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 以感知重试时机。 return raiseRustuse 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 以感知重试时机。 } Err(e) return Err(e.into()), Ok(value) { /* 使用 value */ } }订阅拓扑变更可以把触发器绑定到引擎的发现事件上实时响应系统拓扑变化——这在Worker 恢复上线后继续未完成工作的场景中尤其有用。Trigger触发时机engine::workers-available有 Worker 连接或断开engine::functions-available有 Function 注册或注销这两个事件在引擎源码中是成对定义的常量engine/src/workers/engine_fn/mod.rs 中的TRIGGER_FUNCTIONS_AVAILABLE与TRIGGER_WORKERS_AVAILABLE。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: {}, });Pythonasync 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使用Deserialize JsonSchema的输入结构体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: VecValue } 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 清单Manifestiii.worker.yamliii.worker.yaml是位于 Worker 根目录的清单文件告诉 iii 如何安装依赖、如何运行该 Worker、以及如何透传配置。它同时作用于两类场景iii worker CLI 命令如start、stop、restart见 Using iii / Workers以及 iii 在 config.yaml 中指定后自动拉起的 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一个重要的设计要点清单只描述如何启动 Worker。Worker 一旦运行起来iii 对它一视同仁——无论是由config.yaml拉起的、由iii worker start启动的还是手动运行且使用了 iii SDK 的进程它们与 Engine 的行为完全一致。仓库中引擎自带的内建 Worker 就是一个活例子engine/src/workers/engine_fn/iii.worker.yaml 即为引擎内建函数的清单其结构name / runtime / scripts与上面的示例同构。排障提示如果 Worker 无法正确启动请先检查它的清单并用iii worker logs查看日志。优雅关闭 Worker调用 SDK 的shutdown可以干净地关闭 WebSocket。Engine 会执行一串确定性的清理把该 Worker 的 Functions 与 Triggers 移出注册表、以worker_disconnected事件触发engine::workers-available并用invocation_stopped取消指向它们的在途调用。如果不调用shutdown进程突然退出最终也会达到同样状态——只是要等 Engine 注意到 socket 断开为止优雅关闭让这个过程确定且更快。Node / TypeScript处理 SIGTERMprocess.on(SIGTERM, async () { await worker.shutdown(); process.exit(0); });Pythonimport signal def _on_term(*_): worker.shutdown() raise SystemExit(0) signal.signal(signal.SIGTERM, _on_term)Rust// Rust 线程自身不会让进程保持存活在 main 返回前 // await 此调用让连接线程干净退出。 worker.shutdown_async().await;一次性 / 临时ephemeralWorker是shutdown最重要的应用场景Kubernetes Job、无服务器容器或定时脚本都可以像普通 Worker 一样接入 iii完成工作后调用shutdown()Rust 为shutdown_async().await即可干净离场。小结本文覆盖了 iii Worker 的完整生命周期iii worker init四语言脚手架含幂等重入与--allow-non-empty语义、以III_URL为唯一耦合点的 WebSocket 接入、connecting → connected → available/busy → disconnected状态机、engine::*::list注册表查询、invocation_stopped断连处理与engine::*-available变更订阅、iii.worker.yaml清单以及shutdown优雅关闭。所有关键机制均可在仓库中逐一溯源CLI 参数在 crates/iii-worker/src/cli/init.rs发现触发器常量在 engine/src/workers/engine_fn/mod.rs断连取消语义在 engine/src/invocation/mod.rs。结合 using-iii 目录 中的 Functions / Triggers / Workers 文档即可支撑生产环境下的多 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),仅供参考