ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

iii Python SDK 参考详解:连接 Worker、注册函数与触发器、调用、队列与通道全解

iii Python SDK 参考详解:连接 Worker、注册函数与触发器、调用、队列与通道全解 iii Python SDK 参考详解连接 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 参考文档docs/0-13-0/sdk-reference/python-sdk.mdx.skill.md结合sdk/packages/python下的实际源码系统讲解 iii Python SDK 的公开 API 面从register_worker建立 WebSocket 连接、注册函数与触发器到trigger调用、触发器动作Trigger Action、错误体系、流式通道Channels与日志/遥测。读完后你可以独立完成一个 iii Python worker 的编写、部署与排错并理解 SDK 在断线重连、命名空间、超时等边界条件下的真实行为。1. 安装与包结构安装方式为 pippip install iii-sdk导入名为iii。包元数据见 pyproject.toml其中[project]的name iii-sdk即 pip 包名与导入名不同写requirements.txt时要写iii-sdk写代码时写import iii。SDK 由两个 PyPI 包协作核心客户端位于 iii 包可观测性、HTTP 类型、队列结果类型等基础设施位于 iii_helpers。包根模块__init__.py的__all__明确列出了当前对外导出的符号包括register_worker、InitOptions、TriggerAction、InvocationError、RegistrationRejectedError、IIIConnectionState、EnqueueResult、IStream等可以作为判断哪些名字可以从包根直接导入的权威清单。参考文档中有一处明确的作者注记该页是手工编写的规划中公开面快照planned public surface snapshot最终参考文档将由 SDK 源码生成。因此本文在介绍每个 API 时均以当前仓库源码iii.py、errors.py、iii_constants.py为准核对签名与默认值个别规划名与现行实现有出入的地方主要是错误类型命名会在对应章节中明确指出。2. 连接 workerregister_worker参考文档给出的签名是def register_worker(address: str, options: InitOptions | None None) - IIIaddress是引擎的 SDK WebSocket URLoptions配置 worker 身份与重连行为返回的III实例承载后文所有方法。当前源码中address已放宽为可省略见 register_worker解析顺序为显式传入的address环境变量III_URL由iii compose等监督进程在拉起 worker 时注入同时注入的还有III_NAMESPACE、III_WORKER_NAME常量DEFAULT_ENGINE_URL ws://127.0.0.1:49134iii_constants.py。源码注释特别说明IPv4 回环地址是刻意写死的因为localhost在部分主机上会解析为::1而引擎可能只监听 IPv4。register_worker最多阻塞 30 秒等待连接建立若超时仅记录警告并照常返回客户端后台继续重试连接成功前注册的函数/触发器会被自动补发replay。要观察真实连接状态迁移应使用add_connection_state_listener。2.1InitOptions参数InitOptions是 iii_constants.py 中的 dataclass完整字段与默认值如下文档只说配置 worker 身份与重连源码给出了全部细节字段类型 / 默认值说明worker_namestr \| None默认Noneworker 显示名。非空环境变量III_WORKER_NAME优先于该值再缺省则为hostname:pid。worker_descriptionstr \| None默认None面向人/LLM 的单行描述会出现在engine::workers::list/engine::workers::info的结果中。namespacestr \| None默认Noneworker 所属命名空间缺省回退到III_NAMESPACE环境变量再无则由引擎应用default。命名空间的作用不止于注册worker 及其函数在此注册之后的trigger目标解析、register_trigger绑定也都跟随它除非调用显式指定别的命名空间。enable_metrics_reportingbool默认True通过 OpenTelemetry 上报 worker 指标。invocation_timeout_msint默认30000trigger()调用的默认超时毫秒常量DEFAULT_INVOCATION_TIMEOUT_MS 30000。reconnection_configReconnectionConfig \| None默认None重连行为None时使用DEFAULT_RECONNECTION_CONFIG。otelOtelConfig \| dict \| NoneOpenTelemetry 配置默认启用可传{enabled: False}或设环境变量OTEL_ENABLEDfalse关闭。参考文档提到的SDK 的 OpenTelemetry 接线通过 options 完成即指此字段导出/聚合侧由 iii-observability 负责。headersdict[str, str] \| None附加到 WebSocket 握手头的自定义头。telemetryTelemetryOptions \| None内部 worker 元数据language、project_name、framework、amplitude_api_key。2.2 源码层面的连接机制从 III 类构造与连接逻辑 可以看到几个值得理解的行为后台事件循环每个III实例启动一个专用线程运行独立的asyncio事件循环同步 APItrigger、shutdown等通过run_coroutine_threadsafe投递到该循环。这也解释了为什么文档区分shutdown与shutdown_async——前者给普通线程调用后者给已经处于asyncio上下文的调用方。指数退避重连_reconnect_loop按initial_delay_ms * backoff_multiplier**attempt计算延迟封顶max_delay_ms再叠加 ±jitter_factor的随机抖动max_retries -1表示无限重试达到上限后连接状态置为failed。重连后的恢复流程连接重建后先发送REATTACH帧携带previous_worker_id与reattach_tokentoken 用于证明就是那个 worker同时让引擎提前作废旧连接然后重放全部 trigger type、function、trigger 注册最后冲刷离线期间积压的消息队列队列上限MAX_QUEUE_SIZE 1000超限丢弃最旧。注册冲突的致命性区分引擎推送registrationrejected时FUNCTION_NAMESPACE_CONFLICT只拒收单个函数、worker 继续服务WORKER_NAMESPACE_CONFLICT同名 worker 冲突则直接抛出RegistrationRejectedError并终止不做重连。命名空间的空白值防护InitOptions.namespace显式设为空白字符串会抛ValueError而不是静默回退到default——源码注释解释了这是为了避免整个项目悄悄从错误的命名空间提供服务这类运维上最难发现的问题。一个最小连接示例来自参考文档与 iii-examplefrom iii import register_worker, InitOptions engine_ws_url os.environ.get(III_URL, ws://localhost:49134) worker register_worker( addressengine_ws_url, optionsInitOptions( worker_nameiii-example, otel{enabled: True, service_name: iii-example}, ), )3. 注册函数register_function参考文档的签名def register_function( self, 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: ...request_format/response_format接受 JSON Schema 字典或RegisterFunctionFormat辅助类型注册后会随函数一起存储供 iii console 和 agent 可读的 skills 使用。当前 实现 在此基础上还有几个关键行为schema 自动提取Python 独有request_format/response_format缺省为None时SDK 会从 handler 的类型注解自动推导 JSON Schema入参取第一个参数的类型、出参取返回类型。Node SDK 因 TS 类型在运行时擦除而依赖显式 schemaPython 侧则typed 即有 schema。文档所说的RegisterFunctionFormatPython-only 的声明辅助与直接传 JSON Schema 字典两种形式最终都会汇入同一条注册消息。同步/异步 handler 均可异步 handler 直接在事件循环执行同步 handler 被自动包装到独立线程中运行run_in_executor语义不阻塞事件循环。per-invocation metadata 的签名嗅探只有 handler 显式声明了名为metadata的参数才会收到调用元数据——支持def handler(data, metadataNone)位置参数第二位与def handler(data, *, metadataNone)关键字两种形态其余签名保持原样调用保证既有 handler 完全向后兼容。判定逻辑见_metadata_passing_mode。HttpInvocationConfig第二个参数不是 handler 而是HttpInvocationConfig时注册的是HTTP 调用型函数如 Lambda、Cloudflare Workers引擎不会本地执行它。返回FunctionRef带id属性与unregister()方法可编程地移除函数重复注册同一function_id会抛ValueError。带 Pydantic 模型的注册示例来自源码 docstringfrom 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) # request/response format 从类型注解自动推导 fn.unregister() # 需要时可从引擎注销4. 触发器register_trigger与register_trigger_type4.1 绑定触发器register_trigger参考文档签名def register_trigger(self, trigger: RegisterTriggerInput | dict[str, Any]) - Trigger: ...要点取消触发器只能通过返回值 handle 上的trigger.unregister()完成不存在顶层unregister_trigger方法。当前 实现 的补充细节触发器 ID 由 SDK 自动生成 UUID输入可以是RegisterTriggerInput模型或等价的普通字典type、function_id、可选config、metadata、namespace命名空间默认规则未显式指定namespace时触发器绑定到当前 worker 的命名空间而不是引擎的default——因为函数本来就注册在 worker 的命名空间里触发器若落在default就开了火却解析不到目标。trigger worker.register_trigger({ type: http, function_id: greet, config: {api_path: /greet, http_method: GET}, }) trigger.unregister() # 通过返回的 handle 注销4.2 声明自定义触发器类型register_trigger_type/unregister_trigger_type参考文档签名def register_trigger_type( self, trigger_type: RegisterTriggerTypeInput | dict[str, Any], handler: TriggerHandler[Any], ) - TriggerTypeRef[Any, Any]: ... def unregister_trigger_type( self, trigger_type: RegisterTriggerTypeInput | dict[str, Any], ) - None: ...即register_trigger_type向引擎声明本 worker 提供的一种新触发器类型unregister_trigger_type移除之。实现位于 iii.py配套的抽象基类在 triggers.pyTriggerHandler是抽象基类要求实现register_trigger(config)与unregister_trigger(config)两个异步方法分别在该类型的触发器实例被注册/注销时回调回调收到的TriggerConfig携带id、function_id、config、metadata与namespace。触发器类型输入支持trigger_request_format/call_request_format两个可选 schema传 Pydantic 类或字典SDK 会把 Pydantic 类转换为 JSON Schema 再上线路。返回的TriggerTypeRef是一个类型化的 handle提供register_trigger(function_id, config, metadataNone)与register_function(function_id, handler)两个方法把类型声明—函数注册—触发器实例绑定串成一次类型安全的操作。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))webhook.register_trigger与底层worker.register_trigger的差异在于前者自动把命名空间设为当前 worker 的命名空间函数与触发器必须同域才能解析后者保持引擎默认语义。5. 调用函数trigger/trigger_async参考文档def trigger(self, request: dict[str, Any] | TriggerRequest) - Any: ... async def trigger_async(self, request: dict[str, Any] | TriggerRequest) - Any: ...trigger是同步入口在 SDK 内部事件循环上跑异步机制trigger_async是给asyncio内部调用方的 awaitable 形式。两者接受TriggerRequest或等价的字典function_id、payload、可选action/timeout_ms/metadata/namespace。返回值语义由 action 决定参考文档与 trigger_async 实现 一致无 action同步等待返回函数值TriggerAction.Enqueue(...)经命名队列异步入队返回带messageReceiptId的字典EnqueueResult在包根导出TriggerAction.Void()fire-and-forget不等待响应返回None。超时处理timeout_ms缺省时取InitOptions.invocation_timeout_ms默认 30000msasyncio.wait_for超时后本地抛出 code 为TIMEOUT的调用错误。命名空间解析上以engine::开头的引擎内建函数固定在default命名空间其余隐式调用继承当前 worker 的命名空间显式namespace永远优先。result worker.trigger({function_id: greet, payload: {name: World}}) worker.trigger({function_id: notify, payload: {}, action: TriggerAction.Void()}) receipt worker.trigger({function_id: process, payload: {}, action: TriggerAction.Enqueue(queuejobs)})5.1 Trigger ActionTriggerAction是工厂类提供两个静态方法TriggerActionEnqueue与TriggerActionVoid是具体返回形态定义见 iii.py 与 iii_types.pyTriggerAction.Void() # fire-and-forget; 返回 TriggerActionVoid() TriggerAction.Enqueue(queuemath) # 经命名队列路由; 返回 TriggerActionEnqueue(...)queue在Enqueue上是关键字参数keyword-only即必须写作Enqueue(queue...)。队列名必须在队列 worker 的queue_configs中事先声明。若action以普通字典形式传入{type: enqueue, queue: ...}trigger_async内部会还原为对应的模型对象。6. 错误体系参考文档描述的是规划形态基类IIIInvocationErrorwire 错误解码器把引擎错误码映射到两个已知子类其余留在基类上规划类名触发条件IIIInvocationError任意引擎侧调用错误IIIForbiddenErrorcode FORBIDDENRBAC 拒绝IIITimeoutErrorcode TIMEOUT三者均从包根导出基类携带code、message、function_id、stacktrace属性。与当前源码的对应关系现在 errors.py 中实际导出的基类名为InvocationError属性为code、message、function_id、stacktrace另加invocation_id按错误码分支的编码风格保留了下来——调用方捕获InvocationError后检查err.codeTIMEOUT表示超时由 SDK 在本地超时路径构造FORBIDDEN表示 RBAC 拒绝。wire 错误解码器_wrap_wire_error把引擎错误帧ErrorBody形态的字典转换为InvocationError畸形错误帧回落到codeUNKNOWN保证任何拒绝路径都不会打印原始字典。除调用错误外当前包根还导出RegistrationRejectedErrorerrors.py引擎以WORKER_NAMESPACE_CONFLICT等理由拒绝注册时抛出携带code、namespace、worker_name、owner_worker_id且属于致命错误——SDK 不会重连。stacktrace属性是引擎侧追踪信息可能包含内部文件路径文档与源码一致地提醒不要直接暴露给终端用户str(err)也刻意不含它。from iii import InvocationError, RegistrationRejectedError try: result worker.trigger({function_id: charge, payload: {...}}) except InvocationError as err: if err.code FORBIDDEN: ... # RBAC 拒绝 elif err.code TIMEOUT: ... # 超时7. Channels流式通道参考文档说明ChannelReader与ChannelWriter封装引擎的流式 WebSocketStreamChannelRef标识一个通道class StreamChannelRef(BaseModel): channel_id: str access_key: str direction: Literal[read, write]两个类都以引擎 WS base URL 一个StreamChannelRef构造。当前 channels.py 的实现补充了关键细节通道 URL 由build_channel_url拼出{engine_ws_base}/ws/channels/{channel_id}?key{access_key}dir{direction}其中 access key 经过 URL 编码单帧上限MAX_FRAME_SIZE 64 * 102464 KB超过的负载需要自行分帧ChannelItem表示一帧text与binary二选一text_item/binary_item工厂方法此外还提供 Node 风格的ReadableStream/WritableStream适配层write为 fire-and-forget 二进制写、end可附带最后一帧数据并关闭流。通道引用会随调用负载到达 workerSDK 在_resolve_channels中递归遍历负载凡识别为通道引用的字段按direction就地替换为ChannelReader或ChannelWriter实例字典/列表/元组结构都会被递归处理。iii.helpers.create_channel底层调用引擎内建函数engine::channels::create则是创建通道对writer reader的公共入口。8. Logger 与遥测参考文档Logger暴露info、warn、error、debug每个方法接受一条消息和一个可选的数据字典输出与 SDK 的 OpenTelemetry 配置集成导出侧见 iii-observability。源码层面logger.py每条日志自动捕获当前 trace/span 上下文与分布式追踪自动关联无需手工接线结构化数据以键值字典传入而非字符串拼接便于在可观测后端做过滤、聚合和看板OTel 未初始化时优雅降级到 Python 标准logging日志级别映射到 OTel SeverityNumberdebug→5、info→9、warn→13、error→17。logger.info(Order processed, {order_id: ord_123, amount: 49.99, currency: USD}) logger.warn(Retry attempt, {attempt: 3, max_retries: 5})调用链上的追踪同样内置SDK 在每次调用中创建名为execute function_id的 INTERNAL span刻意不叫call/trigger避免与引擎侧 span 重名并注入traceparent/baggage传播 W3C trace context是否把输入/输出负载打进 span 事件可用III_DISABLE_TRACE_PAYLOADS1|true关闭。9. 连接状态、线协议与其他类型IIIConnectionState字面量类型别名disconnected | connecting | connected | reconnecting | failed定义于 iii_constants.py。参考文档指出它未从包根再导出需从iii.iii_constants导入就当前仓库而言包根__init__.py已将其列入导出清单可直接from iii import IIIConnectionState。文档建议的实践经验仍然成立register_worker返回即可把连接视为已建立即便 30 秒等待超时客户端仍会返回并后台续连。worker.get_connection_state()可随时查询当前状态add_connection_state_listener可订阅状态迁移回调在后台事件循环线程上触发应保持轻量且不要在其中调用同步 SDK 方法。MessageType运行时枚举命名 SDK 与引擎之间交换的每一类线帧如INVOKE_FUNCTION、INVOCATION_RESULT、REATTACH、WORKER_REGISTERED、REGISTRATION_REJECTED等定义于 iii_types.py。它主要供中间件内部使用一般调用方无需接触。Info 类型参考文档列出了FunctionInfofunction_id、可选description、可选request_format/response_format、可选metadata与TriggerInfoid、trigger_type、function_id、可选config/metadataWorkerInfo存在于iii_types但当时未从包根再导出需from iii.iii_types import WorkerInfo而WorkerMetadata不属于本 SDK。需要注意当前源码的 包根导出清单 中已不再包含这些 Info 类型示例工程的注释也确认SDK 不再附带FunctionInfo/FunctionSummary类——即查询函数/触发器清单应走引擎内建函数engine::functions::list等常量见 EngineFunctions而非本地类型。RegisterFunctionFormatPython-only 的辅助类型以结构化方式声明函数入参/出参 schema它与直接传 JSON Schema 字典两种形式都通用于register_function。10. 完整示例仓库内的参考 workeriii-example 是仓库内置的完整 Python worker 演示TODO 应用展示了真实项目的接入姿态engine_ws_url os.environ.get(III_URL, ws://localhost:49134) iii register_worker( addressengine_ws_url, optionsInitOptions( worker_nameiii-example, otel{enabled: True, service_name: iii-example}, ), ) # 之后通过 hooks/streams/state 注册 HTTP 端点、自定义触发器类型、 # 流式存储等进程保持存活以持续服务引擎调用该工程还配有 worker-compose.yaml 与 docker-compose.yaml说明 worker 既可直接以进程方式运行也可纳入 iii 的 compose 编排III_URL、III_WORKER_NAME、III_NAMESPACE等环境契约在 环境变量契约测试 中有对应验证。11. 小结API 面速查能力入口关键行为连接register_worker(addressNone, optionsNone)地址解析参数 →III_URL→ws://127.0.0.1:49134最多阻塞 30 秒断线指数退避重连并重放注册函数register_function(...) - FunctionRefschema 可从类型注解自动推导同步 handler 自动入线程ref.unregister()注销触发器register_trigger(...) - TriggerID 由 SDK 生成 UUIDtrigger.unregister()注销无顶层注销方法触发器类型register_trigger_type(...) - TriggerTypeRef返回类型化 handle串联类型声明/函数/实例绑定调用trigger/trigger_async同步/异步入队/Void 三种语义默认超时 30000ms错误InvocationErrorRegistrationRejectedError按err.code分支TIMEOUT、FORBIDDEN等通道ChannelReader/ChannelWriterStreamChannelRef独立 WS/ws/channels/{id}?key..dir..帧上限 64KB日志Logger.info/warn/error/debug自动关联 trace/spanOTel 未初始化时降级标准 logging参考文档作为规划面快照的价值在于它固定了 SDK 的公共 API 契约哪些方法属于III、哪些名字从包根导出、注销语义放在 handle 上而非顶层而sdk/packages/python源码则提供了契约背后的真实机制重连重放、命名空间冲突的致命性分级、超时与错误的本地构造、负载中通道引用的自动水合。二者对照阅读即可把一个 iii Python 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),仅供参考
返回列表