
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 SDKiii-sdk的完整 API 体系从register_worker连接引擎、register_function注册与类型提示自动抽取 schema、trigger三种路由方式调用函数到自定义触发器类型注册与连接状态管理。读完本文你可以直接在 Python 中编写、注册并运维一个完整的 iii Worker并理解各参数在 SDK 源码 中的真实行为。安装与适用前提pip install iii-sdk根据 pyproject.toml 的定义iii-sdk包要求Python 3.10classifiers 中声明支持 3.10 / 3.11 / 3.12核心依赖为websockets12.0与引擎通信的 WebSocket 传输pydantic2.0类型 schema 抽取Pydantic 模型可直接作为 format 参数opentelemetry-api1.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_workerregister_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_namemy-worker))参数说明addressstr | NoneIII 引擎的 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。optionsInitOptions | NoneWorker 名称、超时、重连与 OTel 等配置完整字段如下源自 InitOptions dataclass字段类型默认说明worker_namestr \| NoneNoneWorker 显示名默认hostname:pid非空的III_WORKER_NAME环境变量会覆盖它worker_descriptionstr \| NoneNone一行人类/LLM 可读的 Worker 摘要会出现在engine::workers::list/engine::workers::info中namespacestr \| NoneNone该 Worker 所属命名空间回退到III_NAMESPACE环境变量两者都未设置时引擎应用default。Worker 及其函数都注册在此后续的trigger目标与register_trigger绑定也默认继承它除非调用方显式指定其他命名空间。对引擎内置engine::*函数的隐式调用解析在default可用显式命名空间覆盖enable_metrics_reportingboolTrue通过 OpenTelemetry 上报 Worker 指标invocation_timeout_msint30000worker.trigger()调用的默认超时毫秒reconnection_configReconnectionConfig \| NoneNoneWebSocket 重连行为缺省时使用DEFAULT_RECONNECTION_CONFIGotelOtelConfig \| dict[str, Any] \| NoneNoneOpenTelemetry 配置默认启用设置{enabled: False}或环境变量OTEL_ENABLEDfalse可关闭headersdict[str, str] \| NoneNone握手头telemetryTelemetryOptions \| NoneNone上报给引擎的内部 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_idstr是函数的唯一字符串标识符handler_or_invocationRemoteFunctionHandler \| HttpInvocationConfig是可调用 handler 或 HTTP 调用配置。handler 第一个参数data接收触发负载可选第二个参数metadata接收每次调用的元数据可返回值descriptionstr \| None否函数用途的人类可读描述metadatadict[str, Any] \| None否附加到函数上的任意元数据request_formatRegisterFunctionFormat \| dict \| None否描述输入期望的 schema为None默认时从 handler 第一个参数的类型提示自动抽取。传显式 schema 可覆盖当 handler 带类型标注时无法注册为无 schemaresponse_formatRegisterFunctionFormat \| dict \| None否描述输出期望的 schema自动抽取语义与request_format相同行为要点与源码一致同步/异步 handler 均支持。同步 handler 会被自动用run_in_executor包装避免阻塞事件循环metadata 只转发给显式声明了参数名为metadata的 handlerdef handler(data, metadataNone)位置参数或def handler(data, *, metadataNone)关键字参数均可判断逻辑见 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: fHello, {data[name]}!} fn worker.register_function(greet, greet, descriptionGreets 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(messagefHello, {data.name}!) fn worker.register_function(greet, greet, descriptionGreets a user)RegisterFunctionFormat的字段结构如下当需要手写 schema 时字段类型必填说明namestr是参数名typestr是类型字符串string、number、boolean、object、array、null、maprequiredbool否是否必填descriptionstr \| None否参数的人类可读描述bodylist[RegisterFunctionFormat] \| None否object 类型的嵌套字段itemsRegisterFunctionFormat \| None否array 类型的元素 schema调用函数trigger / trigger_asynctrigger同步与trigger_async异步调用远程函数。路由行为与返回类型取决于action字段无 action同步请求/响应等待函数返回TriggerAction.Enqueue(queue...)经命名队列异步处理返回含messageReceiptId的 dictEnqueueResultTriggerAction.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_idstr是要调用的函数 IDpayloadAny否传入函数的输入数据actionTriggerActionEnqueue \| TriggerActionVoid \| None否路由方式省略为同步请求/响应namespacestr \| None否路由的目标命名空间省略时继承本 Worker 的命名空间显式写default可从命名空间 Worker 触达引擎默认命名空间metadataAny \| None否每次调用都传递给被触发 handler 的用户自定义元数据须 handler 声明metadata参数timeout_msint \| 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字段类型必填说明typestr是触发器类型标识如storage::object-created、httpfunction_idstr是触发时调用的函数 IDconfigAny否触发器类型专属配置需匹配该触发器类型期望的形状metadataAny \| None否每次调用传递给被触发 handler 的用户自定义元数据namespacestr \| None否目标函数解析所在的命名空间trigger_namespacestr \| None否触发器类型提供者所在的命名空间# dict 形式 trigger worker.register_trigger({ type: http, function_id: greet, config: {api_path: /greet, http_method: GET} }) # RegisterTriggerInput 形式 trigger worker.register_trigger(RegisterTriggerInput( typehttp, function_idgreet, 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字段字段类型必填说明idstr是触发器类型唯一标识如state、durable:subscriberdescriptionstr是触发器类型用途的人类可读描述trigger_request_formatAny \| None否描述期望触发器配置的 JSON SchemaPydantic 类或 dictcall_request_formatAny \| None否描述发送给函数的负载的 JSON Schemahandler须为TriggerHandler实例抽象基类定义于 trigger.py必须实现async register_trigger(config: TriggerConfig[TConfig]) - None按给定配置注册触发器async unregister_trigger(config: TriggerConfig[TConfig]) - None注销触发器。webhook worker.register_trigger_type( RegisterTriggerTypeInput( idwebhook, descriptionWebhook trigger, trigger_request_formatWebhookConfig, call_request_formatWebhookCallRequest, ), WebhookHandler(), ) webhook.register_function(handler, handle_webhook) webhook.register_trigger(handler, WebhookConfig(url/hook))TriggerTypeRef是带两个类型参数Cregister_trigger的配置类型Rregister_function的调用请求类型的类型化句柄register_function(function_id, handler, *, descriptionNone)注册输入与 call-request format 匹配的函数register_trigger(function_id, config, metadataNone) - Trigger以经过校验的配置注册触发器。注销使用unregister_trigger_typeworker.unregister_trigger_type({id: webhook, description: Webhook trigger}) worker.unregister_trigger_type(RegisterTriggerTypeInput(idwebhook, descriptionWebhook 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(fengine 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、failediii_constants.py。get_connection_state / get_addressget_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优雅关闭客户端并释放所有资源。二者语义相同同步/异步版本取消所有挂起的重连尝试以错误拒绝所有在途调用codeSHUTDOWN关闭 WebSocket 连接停止后台事件循环线程。此调用之后实例不可复用。worker register_worker(ws://localhost:49134) # ... do work ... worker.shutdown() # 异步版本 await worker.shutdown_async()核心类型速查错误类型iii.errorsInvocationErrorSDK 派发的调用失败时抛出。检查err.code应对特定类别如 RBAC 拒绝的FORBIDDEN、超时的TIMEOUT捕获该类型可处理所有拒绝。因其继承自Exceptionexcept 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_typesMessageType引擎通信的消息类型常量包括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_REGISTEREDRegisterFunctionInput/RegisterFunctionMessage函数注册输入id必填可选description、invocationHttpInvocationConfig用于外部托管函数、metadata、request_format、response_format及对应线上消息RegisterTriggerInput/RegisterTriggerMessage触发器注册输入与线上消息消息体额外含生成的id与message_typetype字段在消息中为trigger_typeRegisterTriggerTypeInput/RegisterTriggerTypeMessage触发器类型注册输入与线上消息。队列与遥测iiiEnqueueResultTriggerAction.Enqueue调用的返回结果仅含messageReceiptId入队消息的唯一回执 IDTelemetryOptions上报给引擎的 Worker 元数据字段language、project_name、framework、amplitude_api_key。流式通道iii.channel用于 Worker 间数据传输的 WebSocket 流式通道Channel通道对含reader/writerChannelReader/ChannelWriter及reader_ref/writer_refStreamChannelRefChannelReaderread_all() - bytes读完整流、on_message(callback)、close_async()ChannelWriterwrite(data: bytes)、send_message(msg)fire-and-forget向运行中的循环排队协程、send_message_async、close()、close_async()StreamChannelRefchannel_id、access_key认证通道访问的秘密键、directionread/write。流触发器相关类型StreamRequest注册了 stream 触发器的函数收到的流式请求——body、headers、method、path_params、query_params、request_bodyChannelReaderStreamResponse基于ChannelWriter构建——status(status_code)、headers(dict)、close()、writer、stream。引擎常量iii.engineEngineFunctions引擎内置函数 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.pyEngineTriggers引擎触发器 IDengine::functions-available、log。运行时句柄iii.runtimeFunctionRefidunregister()支持程序化反注册TriggerTypeRef上文已述的类型化句柄。状态与流接口iii.state / iii.streamIState状态管理抽象接口按scope命名空间key操作方法签名说明getasync (StateGetInput) - TData \| None按键取值setasync (StateSetInput) - StateSetResult \| None创建或覆盖deleteasync (StateDeleteInput) - StateDeleteResult删除listasync (StateListInput) - list[TData]列出 scope 内全部值updateasync (StateUpdateInput) - StateUpdateResult \| None原子应用list[UpdateOp]更新操作输入/结果类型均为scopekeyset额外含valueupdate含ops的 dataclass结果携带new_value/old_value。StateEventData描述状态变更事件负载event_typeCREATED/UPDATED/DELETED、scope、key、old_value、new_value、type恒为state。IStream流操作抽象接口按stream_namegroup_iditem_id定位条目方法签名说明getasync (StreamGetInput) - TData \| None取单条setasync (StreamSetInput) - StreamSetResult \| None写入data字段deleteasync (StreamDeleteInput) - StreamDeleteResult删除listasync (StreamListInput) - list[TData]列出组内全部条目list_groupsasync (StreamListGroupsInput) - List[str]列出流内全部组updateasync (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),仅供参考