iii Python SDK 完整参考:从 register_worker 注册 Worker 到函数调用与自定义触发器类型
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
本文系统讲解 iii 项目 Python SDK(iii-sdk)的完整 API 体系:从register_worker连接引擎、register_function注册与类型提示自动抽取 schema、trigger三种路由方式调用函数,到自定义触发器类型注册与连接状态管理。读完本文,你可以直接在 Python 中编写、注册并运维一个完整的 iii Worker,并理解各参数在 SDK 源码 中的真实行为。
安装与适用前提
pip install iii-sdk根据 pyproject.toml 的定义,iii-sdk包要求Python >= 3.10(classifiers 中声明支持 3.10 / 3.11 / 3.12),核心依赖为:
websockets>=12.0:与引擎通信的 WebSocket 传输;pydantic>=2.0:类型 schema 抽取(Pydantic 模型可直接作为 format 参数);opentelemetry-api>=1.25:可观测性支持;iii-helpers:同仓库 sdk/packages/python/helpers 下的本地依赖包(含HttpInvocationConfig、OtelConfig、ReconnectionConfig、EnqueueResult等共享类型)。
SDK 的公共入口在 iii/init.py 中统一导出:register_worker、InitOptions、TriggerAction、InvocationError、RegistrationRejectedError、EnqueueResult、IIIClient、IStream等。
本文对应的参考文档由源码 doc-comment 自动渲染生成(见 sdk-python.mdx.skill.md 头部注释),修改点位于
sdk/packages/python/iii/src下的源文档注释;所有行为描述均可在该目录的源码与 tests 目录 的测试用例中得到印证。
初始化:register_worker
register_worker是 SDK 的主入口,创建 Worker 客户端并将其注册到 iii 引擎,返回一个已连接的III客户端。
签名
register_worker(address: str | None = None, options: InitOptions | None = None) -> IIIfrom iii import register_worker, InitOptions worker = register_worker() # address 从 III_URL 解析 other = register_worker('ws://localhost:49134', InitOptions(worker_name='my-worker'))参数说明
address(str | None):III 引擎的 WebSocket URL(如ws://localhost:49134)。省略时按顺序从环境变量III_URL解析,最终回退到DEFAULT_ENGINE_URL。在 iii_constants.py 中,该默认值为ws://127.0.0.1:49134——源码注释特别说明这里刻意写死 IPv4 loopback,因为localhost在部分主机上会解析到::1,而引擎可能只监听 IPv4。解析逻辑见 iii.py 的resolve_engine_url:显式参数 →III_URL→DEFAULT_ENGINE_URL。options(InitOptions | None):Worker 名称、超时、重连与 OTel 等配置,完整字段如下(源自 InitOptions dataclass):
| 字段 | 类型 | 默认 | 说明 |
|---|---|---|---|
worker_name | str \| None | None | Worker 显示名,默认hostname:pid;非空的III_WORKER_NAME环境变量会覆盖它 |
worker_description | str \| None | None | 一行人类/LLM 可读的 Worker 摘要,会出现在engine::workers::list/engine::workers::info中 |
namespace | str \| None | None | 该 Worker 所属命名空间,回退到III_NAMESPACE环境变量;两者都未设置时引擎应用default。Worker 及其函数都注册在此,后续的trigger目标与register_trigger绑定也默认继承它(除非调用方显式指定其他命名空间)。对引擎内置engine::*函数的隐式调用解析在default,可用显式命名空间覆盖 |
enable_metrics_reporting | bool | True | 通过 OpenTelemetry 上报 Worker 指标 |
invocation_timeout_ms | int | 30000 | worker.trigger()调用的默认超时(毫秒) |
reconnection_config | ReconnectionConfig \| None | None | WebSocket 重连行为,缺省时使用DEFAULT_RECONNECTION_CONFIG |
otel | OtelConfig \| dict[str, Any] \| None | None | OpenTelemetry 配置,默认启用;设置{'enabled': False}或环境变量OTEL_ENABLED=false可关闭 |
headers | dict[str, str] \| None | None | 握手头 |
telemetry | TelemetryOptions \| None | None | 上报给引擎的内部 Worker 元数据 |
从源码结构看(III 构造函数),客户端在构造时即启动一个后台事件循环线程并自动发起连接;_wait_until_connected最多阻塞30 秒等待 WebSocket 建立。若超时,仅记录 warning 并照常返回客户端——它会继续在后台重试,连接恢复后统一 flush 已排队的注册消息。要观察真实的连接状态迁移,应使用下文add_connection_state_listener。
命名空间解析有一些值得注意的细节(III 类中的_call_namespace/_worker_namespace):
- 显式传入空字符串/纯空白命名空间会直接抛出
ValueError,而不是静默回退——源码注释指出 Python 中""为 falsy,若沿用or逻辑会导致"未设置"与"设置了空名"两种相反语义被混淆; - 环境变量
III_NAMESPACE若为空白则按"未设置"处理(与 shell 中III_NAMESPACE=${NS}展开为空的习惯一致); - 对以
engine::开头的函数调用,未显式指定命名空间时自动解析到default(见 _invocation_namespace)。
注册函数:register_function
将函数注册到引擎。可以传入本地执行的 handler,也可以传入HttpInvocationConfig用于 HTTP 调用的外部函数(Lambda、Cloudflare Workers 等)。
签名
register_function( function_id: str, handler_or_invocation: RemoteFunctionHandler | HttpInvocationConfig, *, description: str | None = None, metadata: dict[str, Any] | None = None, request_format: RegisterFunctionFormat | dict[str, Any] | None = None, response_format: RegisterFunctionFormat | dict[str, Any] | None = None, ) -> FunctionRef参数说明
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
function_id | str | 是 | 函数的唯一字符串标识符 |
handler_or_invocation | RemoteFunctionHandler \| HttpInvocationConfig | 是 | 可调用 handler 或 HTTP 调用配置。handler 第一个参数data接收触发负载,可选第二个参数metadata接收每次调用的元数据,可返回值 |
description | str \| None | 否 | 函数用途的人类可读描述 |
metadata | dict[str, Any] \| None | 否 | 附加到函数上的任意元数据 |
request_format | RegisterFunctionFormat \| dict \| None | 否 | 描述输入期望的 schema;为None(默认)时从 handler 第一个参数的类型提示自动抽取。传显式 schema 可覆盖;当 handler 带类型标注时,无法注册为"无 schema" |
response_format | RegisterFunctionFormat \| dict \| None | 否 | 描述输出期望的 schema,自动抽取语义与request_format相同 |
行为要点(与源码一致)
- 同步/异步 handler 均支持。同步 handler 会被自动用
run_in_executor包装,避免阻塞事件循环; - metadata 只转发给显式声明了参数名为
metadata的 handler:def handler(data, metadata=None)(位置参数)或def handler(data, *, metadata=None)(关键字参数)均可,判断逻辑见 iii.py 的_metadata_passing_mode。签名无法内省(部分 builtin/C 可调用对象)时回退为"不转发",因此既有 handler——包括带其他可选参数、*args、**kwargs的——保持原样工作; - schema 自动抽取是 Python 特有行为:
request_format/response_format在省略或传None时从类型提示自动提取(Pydantic 模型经python_type_to_format转换为 JSON Schema,见_resolve_format)。Node SDK 因 TypeScript 类型在运行时被擦除,依赖显式 schema; - 返回的
FunctionRef提供.unregister()用于程序化反注册。
示例
# 简单 dict handler def greet(data): return {'message': f"Hello, {data['name']}!"} fn = worker.register_function("greet", greet, description="Greets a user") fn.unregister() # 使用 Pydantic 模型,request/response format 自动抽取 from pydantic import BaseModel class GreetInput(BaseModel): name: str class GreetOutput(BaseModel): message: str async def greet(data: GreetInput) -> GreetOutput: return GreetOutput(message=f"Hello, {data.name}!") fn = worker.register_function("greet", greet, description="Greets a user")RegisterFunctionFormat的字段结构如下(当需要手写 schema 时):
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
name | str | 是 | 参数名 |
type | str | 是 | 类型字符串:string、number、boolean、object、array、null、map |
required | bool | 否 | 是否必填 |
description | str \| None | 否 | 参数的人类可读描述 |
body | list[RegisterFunctionFormat] \| None | 否 | object 类型的嵌套字段 |
items | RegisterFunctionFormat \| None | 否 | array 类型的元素 schema |
调用函数:trigger / trigger_async
trigger(同步)与trigger_async(异步)调用远程函数。路由行为与返回类型取决于action字段:
- 无 action:同步请求/响应,等待函数返回;
TriggerAction.Enqueue(queue=...):经命名队列异步处理,返回含messageReceiptId的 dict(EnqueueResult);TriggerAction.Void():fire-and-forget,返回None。
签名
trigger(request: dict[str, Any] | TriggerRequest) -> Any async trigger_async(request: dict[str, Any] | TriggerRequest) -> Anyresult = worker.trigger({'function_id': 'greet', 'payload': {'name': 'World'}}) worker.trigger({'function_id': 'notify', 'payload': {}, 'action': TriggerAction.Void()}) # 异步版本 result = await worker.trigger_async({'function_id': 'greet', 'payload': {'name': 'World'}}) await worker.trigger_async({'function_id': 'notify', 'payload': {}, 'action': TriggerAction.Void()})TriggerRequest字段说明:
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
function_id | str | 是 | 要调用的函数 ID |
payload | Any | 否 | 传入函数的输入数据 |
action | TriggerActionEnqueue \| TriggerActionVoid \| None | 否 | 路由方式;省略为同步请求/响应 |
namespace | str \| None | 否 | 路由的目标命名空间;省略时继承本 Worker 的命名空间,显式写default可从命名空间 Worker 触达引擎默认命名空间 |
metadata | Any \| None | 否 | 每次调用都传递给被触发 handler 的用户自定义元数据(须 handler 声明metadata参数) |
timeout_ms | int \| None | 否 | 覆盖默认调用超时(毫秒),默认值由InitOptions.invocation_timeout_ms决定(30000) |
关于Enqueue:它需要worker-compose.yaml中存在queueWorker 且其queue_configs有对应条目,否则触发会以enqueue_error(无队列提供者)被拒绝。TriggerActionEnqueue结构为{queue: str, type: Literal['enqueue']},TriggerActionVoid为{type: Literal['void']}。
注册触发器:register_trigger
将触发器配置绑定到一个已注册的函数。
签名
register_trigger(trigger: RegisterTriggerInput | dict[str, Any]) -> Trigger| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
type | str | 是 | 触发器类型标识(如storage::object-created、http) |
function_id | str | 是 | 触发时调用的函数 ID |
config | Any | 否 | 触发器类型专属配置,需匹配该触发器类型期望的形状 |
metadata | Any \| None | 否 | 每次调用传递给被触发 handler 的用户自定义元数据 |
namespace | str \| None | 否 | 目标函数解析所在的命名空间 |
trigger_namespace | str \| None | 否 | 触发器类型提供者所在的命名空间 |
# dict 形式 trigger = worker.register_trigger({ 'type': 'http', 'function_id': 'greet', 'config': {'api_path': '/greet', 'http_method': 'GET'} }) # RegisterTriggerInput 形式 trigger = worker.register_trigger(RegisterTriggerInput( type="http", function_id="greet", config={'api_path': '/greet', 'http_method': 'GET'} )) trigger.unregister()返回的Trigger句柄提供.unregister()方法。
自定义触发器类型:register_trigger_type / unregister_trigger_type
将自定义触发器类型注册到引擎,返回带register_trigger与register_function方法的TriggerTypeRef句柄。
签名
register_trigger_type( trigger_type: RegisterTriggerTypeInput | dict[str, Any], handler: TriggerHandler[Any], ) -> TriggerTypeRef[Any, Any]RegisterTriggerTypeInput字段:
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
id | str | 是 | 触发器类型唯一标识(如state、durable:subscriber) |
description | str | 是 | 触发器类型用途的人类可读描述 |
trigger_request_format | Any \| None | 否 | 描述期望触发器配置的 JSON Schema(Pydantic 类或 dict) |
call_request_format | Any \| None | 否 | 描述发送给函数的负载的 JSON Schema |
handler须为TriggerHandler实例(抽象基类,定义于 trigger.py),必须实现:
async register_trigger(config: TriggerConfig[TConfig]) -> None:按给定配置注册触发器;async unregister_trigger(config: TriggerConfig[TConfig]) -> None:注销触发器。
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"))TriggerTypeRef是带两个类型参数(C:register_trigger的配置类型;R:register_function的调用请求类型)的类型化句柄:
register_function(function_id, handler, *, description=None):注册输入与 call-request format 匹配的函数;register_trigger(function_id, config, metadata=None) -> Trigger:以经过校验的配置注册触发器。
注销使用unregister_trigger_type:
worker.unregister_trigger_type({"id": "webhook", "description": "Webhook trigger"}) worker.unregister_trigger_type(RegisterTriggerTypeInput(id="webhook", description="Webhook trigger"))TriggerConfig是注册/注销时传给 handler 的配置对象:id(触发器实例 ID)、function_id、config(触发器专属配置)、metadata、namespace(当前 SDK 会用注册 Worker 的命名空间填充省略值;None为遗留/默认情形)。
连接状态与生命周期管理
add_connection_state_listener
订阅连接状态迁移事件。
unsubscribe = worker.add_connection_state_listener( lambda state: print(f"engine link: {state}") )add_connection_state_listener(handler: ConnectionStateCallback) -> Callable[[], None]行为约定(与 III 实现 一致):
- handler 会立即以当前状态触发一次(在调用方线程上),之后每次状态迁移再触发;
- 迁移回调在 SDK 后台事件循环线程上执行:handler 要保持轻量,不要在 handler 中调用同步 SDK 方法(会抛出
RuntimeError,见_run_on_loop的线程检查); - 将回调视为"状态通知"而非"状态边沿",某状态在订阅前后可能被观察到两次;
- 同一 handler 注册两次会触发两次;返回的 unsubscribe 函数是幂等的,且只移除自己的注册。
连接状态IIIConnectionState的取值为:disconnected、connecting、connected、reconnecting、failed(iii_constants.py)。
get_connection_state / get_address
get_connection_state() -> IIIConnectionStateworker = register_worker("ws://localhost:49134") if worker.get_connection_state() != "connected": print("engine not reachable yet")get_address() -> str返回该 Worker 实际解析到的引擎地址:显式register_worker参数 →III_URL→DEFAULT_ENGINE_URL。与 Rust SDK 的address()和 Node SDK 的getAddress()对齐。
connect_async
通过 WebSocket 连接 III 引擎:初始化 OpenTelemetry(如已配置)、附加事件循环、建立 WebSocket 连接。该调用在构造时已自动执行,仅在需要从异步上下文手动重连时使用。
async connect_async() -> None从 源码 可见其内部顺序:init_otel→attach_event_loop→ 状态置connecting→_do_connect。
shutdown / shutdown_async
优雅关闭客户端并释放所有资源。二者语义相同(同步/异步版本):
- 取消所有挂起的重连尝试;
- 以错误拒绝所有在途调用(
code="SHUTDOWN"); - 关闭 WebSocket 连接;
- 停止后台事件循环线程。
此调用之后实例不可复用。
worker = register_worker('ws://localhost:49134') # ... do work ... worker.shutdown() # 异步版本 await worker.shutdown_async()核心类型速查
错误类型(iii.errors)
InvocationError:SDK 派发的调用失败时抛出。检查err.code应对特定类别(如 RBAC 拒绝的'FORBIDDEN'、超时的'TIMEOUT');捕获该类型可处理所有拒绝。因其继承自Exception,except Exception仍然有效。属性构造后只读;stacktrace是远端 handler 抛出时的引擎侧堆栈,可能包含内部文件路径,不应暴露给终端用户,str(err)也刻意不含堆栈。字段:code、function_id、invocation_id、message、stacktrace。RegistrationRejectedError:引擎拒绝本 Worker 的注册时抛出。注册冲突(如另一个存活 Worker 已拥有(namespace, worker_name))时,引擎推送registrationrejected消息并关闭连接。这是致命的:SDK 不会重连。字段:code、namespace、owner_worker_id、worker_name。
消息与协议类型(iii.protocol / iii.iii_types)
MessageType:引擎通信的消息类型常量,包括INVOKE_FUNCTION、INVOCATION_RESULT、REGISTER_FUNCTION、REGISTER_TRIGGER、REGISTER_TRIGGER_TYPE、UNREGISTER_FUNCTION、UNREGISTER_TRIGGER、UNREGISTER_TRIGGER_TYPE、REGISTER_SERVICE、REATTACH、REGISTRATION_REJECTED、TRIGGER_REGISTRATION_RESULT、WORKER_REGISTERED;RegisterFunctionInput/RegisterFunctionMessage:函数注册输入(id必填,可选description、invocation(HttpInvocationConfig,用于外部托管函数)、metadata、request_format、response_format)及对应线上消息;RegisterTriggerInput/RegisterTriggerMessage:触发器注册输入与线上消息(消息体额外含生成的id与message_type,type字段在消息中为trigger_type);RegisterTriggerTypeInput/RegisterTriggerTypeMessage:触发器类型注册输入与线上消息。
队列与遥测(iii)
EnqueueResult:TriggerAction.Enqueue调用的返回结果,仅含messageReceiptId(入队消息的唯一回执 ID);TelemetryOptions:上报给引擎的 Worker 元数据,字段language、project_name、framework、amplitude_api_key。
流式通道(iii.channel)
用于 Worker 间数据传输的 WebSocket 流式通道:
Channel:通道对,含reader/writer(ChannelReader/ChannelWriter)及reader_ref/writer_ref(StreamChannelRef);ChannelReader:read_all() -> bytes(读完整流)、on_message(callback)、close_async();ChannelWriter:write(data: bytes)、send_message(msg)(fire-and-forget,向运行中的循环排队协程)、send_message_async、close()、close_async();StreamChannelRef:channel_id、access_key(认证通道访问的秘密键)、direction(read/write)。
流触发器相关类型:
StreamRequest:注册了 stream 触发器的函数收到的流式请求——body、headers、method、path_params、query_params、request_body(ChannelReader);StreamResponse:基于ChannelWriter构建——status(status_code)、headers(dict)、close()、writer、stream。
引擎常量(iii.engine)
EngineFunctions:引擎内置函数 ID(与 Node SDK 对齐):engine::functions::list/info、engine::workers::list/info、engine::triggers::list/info、engine::registered-triggers::list/info、engine::workers::register(常量定义见 iii_constants.py);EngineTriggers:引擎触发器 ID:engine::functions-available、log。
运行时句柄(iii.runtime)
FunctionRef:id+unregister(),支持程序化反注册;TriggerTypeRef:上文已述的类型化句柄。
状态与流接口(iii.state / iii.stream)
IState:状态管理抽象接口,按scope(命名空间)+key操作:
| 方法 | 签名 | 说明 |
|---|---|---|
get | async (StateGetInput) -> TData \| None | 按键取值 |
set | async (StateSetInput) -> StateSetResult \| None | 创建或覆盖 |
delete | async (StateDeleteInput) -> StateDeleteResult | 删除 |
list | async (StateListInput) -> list[TData] | 列出 scope 内全部值 |
update | async (StateUpdateInput) -> StateUpdateResult \| None | 原子应用list[UpdateOp]更新操作 |
输入/结果类型均为scope+key(set额外含value,update含ops)的 dataclass;结果携带new_value/old_value。StateEventData描述状态变更事件负载:event_type(CREATED/UPDATED/DELETED)、scope、key、old_value、new_value、type(恒为state)。
IStream:流操作抽象接口,按stream_name+group_id+item_id定位条目:
| 方法 | 签名 | 说明 |
|---|---|---|
get | async (StreamGetInput) -> TData \| None | 取单条 |
set | async (StreamSetInput) -> StreamSetResult \| None | 写入(data字段) |
delete | async (StreamDeleteInput) -> StreamDeleteResult | 删除 |
list | async (StreamListInput) -> list[TData] | 列出组内全部条目 |
list_groups | async (StreamListGroupsInput) -> List[str] | 列出流内全部组 |
update | async (StreamUpdateInput) -> StreamUpdateResult \| None | 原子应用list[UpdateOp] |
StreamUpdateResult额外包含errors字段:merge/append对校验拒绝(路径深度/大小、值深度、__proto__/constructor/prototype段或顶级键)及append.type_mismatch、append.target_not_object会逐 op 报错;成功应用的 op 仍反映在new_value中;该字段为空时不出现在 JSON 线上。
总结
iii Python SDK 的 API 设计围绕"注册即连接"的模型展开:register_worker一行代码完成引擎地址解析、后台事件循环启动与 WebSocket 连接,随后用register_function(支持 Pydantic 类型提示自动 schema 抽取与 metadata 参数探测)、trigger(同步 / 队列 / fire-and-forget 三种路由)、register_trigger与register_trigger_type构成完整的 Worker 开发闭环,配合连接状态监听器、EnqueueResult回执与类型化错误(InvocationError/RegistrationRejectedError)实现可观测、可恢复的运行时。实现细节与测试可进一步参阅 sdk/packages/python/iii/src/iii 源码目录及 tests 中的test_sync_api.py、test_trigger_action.py、test_register_function_args.py、test_connection_state_listener.py、test_trigger_type_lifecycle.py等用例。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考