iii-sdk Python 快速上手:用 iii 引擎注册函数、绑定触发器与调用工作流
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
导读
本文以 sdk/packages/python/iii/README.md 为主体,系统讲解 iii-sdk(PyPI 包名iii-sdk)的安装、Worker 初始化、函数注册、触发器绑定与函数调用全流程。iii-sdk 是 iii 引擎的 Python 客户端:Worker 通过 WebSocket 接入引擎,注册可在引擎侧按名称调用的函数,并将 HTTP、cron、队列等触发器绑定到这些函数上。读完本文,你将掌握用 Python 编写可被引擎调度、可被其他 Worker 调用、可接入实时触发器的完整 Worker 的最小实现,以及同步/异步 API、命名空间、连接管理与开发测试配套工具链。
安装与项目环境
iii-sdk 是一个发布在 PyPI 上的标准 Python 包,通过 pip 即可安装:
pip install iii-sdk从 pyproject.toml 可以看到包的工程约束:要求 Python>=3.10,核心依赖为websockets>=12.0(WebSocket 通信)、pydantic>=2.0(消息与配置模型)、opentelemetry-api>=1.25(可观测性)以及同版本的iii-helpers(公共辅助类型,如 HTTP 调用配置与 OTel 配置)。包名与许可证为 Apache-2.0,源码位于 src/iii 目录,顶层导出集中在init.py。
引擎地址的解析规则见 iii.py 中的resolve_engine_url:显式传入的 address 参数优先,其次是环境变量III_URL,最后回退到内置默认值ws://127.0.0.1:49134(见 iii_constants.py,刻意使用 IPv4 回环地址以避免localhost在部分主机上解析为::1的问题)。因此本地开发时,register_worker()可以不传任何参数直接连接默认引擎。
Hello World:完整的最小 Worker
README 给出的 Hello World 是理解 SDK 工作流的最佳入口:
from iii import register_worker iii = register_worker("ws://localhost:49134") def greet(data): return {"message": f"Hello, {data['name']}!"} iii.register_function("hello::greet", greet) iii.register_trigger({ "type": "http", "function_id": "hello::greet", "config": {"api_path": "/greet", "http_method": "POST"}, }) iii.connect() result = iii.trigger({"function_id": "hello::greet", "payload": {"name": "world"}}) print(result) # {"message": "Hello, world!"}这段代码体现了 SDK 的四个核心步骤:
- 初始化连接:
register_worker(url)创建 SDK 实例并自动连接引擎; - 注册函数:
iii.register_function(id, handler)注册一个可被按名称调用的函数; - 绑定触发器:
iii.register_trigger({...})把 HTTP、cron、队列等触发器绑定到函数上; - 调用与关闭:
iii.trigger(...)同步等待调用结果,iii.shutdown()优雅断开连接。
从源码看,register_worker(iii.py)会创建III客户端实例并阻塞等待连接建立(最多 30 秒);若超时未连上,只记录警告并返回客户端,连接会在后台持续重试,注册的消息会先进入发送队列、连接成功后再统一冲刷(见_on_connected与_queue逻辑)。III构造时启动一个非守护后台线程运行独立事件循环,测试 test_sync_api.py 专门验证了这一行为。
API 总览
README 用一张表概括了 SDK 的核心操作,下表在此基础上补充了对应源码位置与关键行为说明:
| 操作 | 签名 | 说明 |
|---|---|---|
| 初始化 | register_worker(url, options?) | 创建 SDK 实例并自动连接引擎;address 可省略,依次从III_URL、默认地址解析 |
| 注册函数 | iii.register_function(id, handler) | 注册可被按名称调用的函数,返回带unregister()的FunctionRef |
| 注册触发器 | iii.register_trigger({"type": ..., "function_id": ..., "config": ...}) | 将 HTTP、cron、队列等触发器绑定到函数,返回带unregister()的Trigger |
| 同步调用(等待结果) | iii.trigger({"function_id": id, "payload": data}) | 发送调用并等待函数返回结果 |
| 异步调用(fire-and-forget) | iii.trigger({..., "action": TriggerAction.Void()}) | 只发送不等待响应,返回None |
| 队列路由调用 | iii.trigger({..., "action": TriggerAction.Enqueue(queue="name")}) | 通过命名队列路由调用,返回包含messageReceiptId的字典 |
| 关闭 | iii.shutdown() | 断开连接并停止后台线程 |
此外,register_worker()的options参数接受InitOptions(定义于 iii_constants.py),常用字段包括:worker_name(Worker 显示名,III_WORKER_NAME环境变量可覆盖,缺省为hostname:pid)、worker_description(一行人类/LLM 可读的职责描述,会出现在engine::workers::list等引擎内建函数中)、namespace(Worker 所属命名空间,回退到III_NAMESPACE环境变量,再缺省由引擎使用default)、invocation_timeout_ms(调用超时,默认 30000)、reconnection_config(重连策略)与otel(OpenTelemetry 配置)。
注册函数:同步、异步与 HTTP 转发
README 的示例展示了返回 HTTP 风格响应的函数注册方式:
def create_order(data): return {"status_code": 201, "body": {"id": "123", "item": data["body"]["item"]}} iii.register_function("orders::create", create_order)从 register_function 的实现 可以看到几处值得注意的行为:
- 同步与异步处理器均可:协程处理器(
asyncio.iscoroutinefunction)直接await;同步处理器会被包装到独立线程中执行(run_in_executor的等价实现),避免阻塞 SDK 的事件循环; - 按需透传调用元数据:只有处理器显式声明名为
metadata的参数时,SDK 才会把每次调用的元数据传给它,支持def handler(data, metadata=None)或def handler(data, *, metadata=None)两种签名(见_metadata_passing_mode,iii.py),既有的一参处理器无需改动; - 请求/响应格式自动提取:省略
request_format/response_format时,SDK 会从处理器的类型注解自动提取 schema(Pydantic 模型会被转换为 JSON Schema);Node SDK 因 TypeScript 类型在运行时被擦除而必须显式传 schema,这是 Python SDK 的独有便利; - HTTP 外部函数:除了本地可调用处理器,还可以传入
iii_helpers.http.HttpInvocationConfig,把函数注册为 HTTP 调用的远端函数(如 Lambda、Cloudflare Workers),引擎会通过 HTTP 转发调用; - 重复注册校验:
function_id为空或已注册会抛出ValueError,非字符串会抛出TypeError; - 返回句柄可注销:
register_function返回的FunctionRef带有unregister(),可编程移除函数(对应引擎侧UnregisterFunctionMessage消息)。
注册触发器:把事件源接到函数上
README 展示了 HTTP 触发器的注册方式:
iii.register_trigger({ "type": "http", "function_id": "orders::create", "config": {"api_path": "/orders", "http_method": "POST"}, })触发器描述符由三要素构成:type(触发器类型,如http、cron、queue等)、function_id(触发时调用的目标函数)、config(类型相关的配置,HTTP 类型下为api_path与http_method)。
源码实现(iii.py)中,SDK 会为每个触发器自动生成 UUID 作为触发器 ID,发送RegisterTriggerMessage给引擎,并返回带unregister()方法的Trigger对象。值得注意的是命名空间语义:触发器默认解析到当前 Worker 的命名空间(而非引擎的default),因为触发器指向的函数注册在 Worker 自己的命名空间里,默认到别处会导致"触发器触发了却解析不到函数"的问题;需要明确指向其他命名空间时,可通过RegisterTriggerInput.namespace显式声明。
自定义触发器类型
除了内置类型,SDK 还支持注册自定义触发器类型。通过register_trigger_type传入类型定义(id、description、可选的trigger_request_format/call_request_formatPydantic 模型)与实现了TriggerHandler抽象基类的处理器(需实现register_trigger/unregister_trigger两个抽象方法,见 triggers.py),即可获得带类型约束的TriggerTypeRef句柄:
webhook = worker.register_trigger_type( RegisterTriggerTypeInput( id="webhook", description="Webhook trigger", trigger_request_format=WebhookConfig, call_request_format=WebhookCallRequest, ), WebhookHandler(), ) webhook.register_function("handler", handle_webhook) webhook.register_trigger("handler", WebhookConfig(url="/hook"))调用函数:同步、Void 与队列三种路由方式
README 中同步调用的写法为:
result = iii.trigger({"function_id": "orders::create", "payload": {"body": {"item": "widget"}}})trigger的完整行为取决于请求中的action字段(源码见 iii.py):
- 不带 action:同步调用,SDK 生成
invocation_id并等待引擎返回结果,支持通过timeout_ms(缺省取InitOptions.invocation_timeout_ms,默认 30000ms)控制超时,超时抛出InvocationError(code="TIMEOUT"); TriggerAction.Void():fire-and-forget,只发送不等待,立即返回None;TriggerAction.Enqueue(queue="name"):通过命名队列路由调用,返回包含messageReceiptId的字典,队列需在队列 Worker 的queue_configs中预先声明。
worker.trigger({'function_id': 'process', 'payload': {}, 'action': TriggerAction.Enqueue(queue='jobs')}) worker.trigger({'function_id': 'notify', 'payload': {}, 'action': TriggerAction.Void()})调用失败时抛出InvocationError,可通过其code字段区分原因:TIMEOUT表示超时,FORBIDDEN表示 RBAC 拒绝。所有同步方法均有对应的trigger_async异步版本(III内部通过asyncio.run_coroutine_threadsafe把同步调用桥接到后台事件循环,iii.py),因此既能在普通脚本中同步调用,也能在 asyncio 程序中直接await worker.trigger_async(...),配套测试见 test_async_api.py。
连接生命周期:重连、状态监听与命名空间
iii-sdk 的连接管理并非"一次连接、断开即死"的简单模型,其关键机制包括:
- 自动重连:连接失败或异常断开后,SDK 按
ReconnectionConfig(initial_delay_ms、backoff_multiplier、max_delay_ms、jitter_factor、max_retries,max_retries=-1表示无限重试)执行带指数退避与抖动的重连循环(_reconnect_loop,iii.py); - 断线重注册:重连成功后,
_on_connected会把此前注册的触发器类型、函数、触发器全部重放给引擎,并冲刷连接期间积压在队列中的消息,保证 Worker 断线期间发起的调用不丢失; - Reattach 身份重续:若已有
worker_id,重连时会先发送REATTACH消息(携带上轮的reattach_token作为身份凭证),让引擎退役旧连接、避免身份冲突; - 状态机与监听:连接状态为
disconnected/connecting/connected/reconnecting/failed五态(iii_constants.py),通过add_connection_state_listener订阅状态迁移(返回幂等退订函数),get_connection_state()随时可查; - 致命拒绝不再重连:若引擎以
REGISTRATION_REJECTED拒绝注册(如WORKER_NAMESPACE_CONFLICT表示同名 Worker 已在存活),SDK 记录致命错误、立即失败所有在途调用并停止重连;而FUNCTION_NAMESPACE_CONFLICT(单个函数 ID 被他人占用)只跳过该函数、不影响 Worker 其余服务; - 命名空间解析:
namespace的优先级为InitOptions.namespace→III_NAMESPACE环境变量 → 引擎default;声明为空白字符串会被视为错误抛出,因为"未设置"与"设为空"语义相反(见 iii.py)。
这些环境变量(III_URL、III_NAMESPACE、III_WORKER_NAME)通常由 iii 的编排器(iii compose、容器运行时或 systemd)注入,让同一份 Worker 代码在不同部署环境下无需改动即可接入不同引擎。
开发、类型检查与测试
README 给出包自身的开发工作流,适用于对 SDK 本身做二次开发或本地调试的场景:
# 开发模式安装 pip install -e . # 类型检查 mypy src # Lint ruff check src工程在 pyproject.toml 中配置了mypy(strict 模式)与ruff(选择 E/F/I/W 规则组,行宽 120),测试采用 pytest 且默认开启分支覆盖率统计(--cov=src/iii --cov-branch)。仓库 tests 目录下有超过 40 个测试文件,覆盖同步/异步 API、触发器注册、重连(如 test_reconnect_sends_reattach.py)、命名空间、RBAC Worker、OpenTelemetry 遥测、流(Streams)、通道(Channels)与状态管理等功能,是理解 SDK 行为契约的补充资料。
总结
iii-sdk(Python)以register_worker→register_function→register_trigger→trigger四步为核心工作流,通过 WebSocket 把 Python Worker 接入 iii 引擎:引擎负责函数路由、触发器分发与跨 Worker 编排,SDK 负责连接管理、消息协议与同步/异步适配。若要进一步深入,可阅读 docs/quickstart.mdx 了解引擎安装与启动方式,参考 sdk/packages/python/iii/README.md 同目录下的源码与测试,或对照 Node、Rust、Go 等其他语言的 SDK 实现理解跨语言一致性设计。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考