ARTICLE DETAIL

资讯详情

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

iii 项目 Python SDK 完整指南:注册函数、触发与流式通信

iii 项目 Python SDK 完整指南:注册函数、触发与流式通信 iii 项目 Python SDK 完整指南注册函数、触发与流式通信【免费下载链接】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代码中导入为iii的权威使用参考。它覆盖从安装、连接引擎register_worker到注册函数register_function、绑定触发器register_trigger/register_trigger_type、发起调用trigger/trigger_async、错误分类、流式 Channel 通信与 OpenTelemetry 日志接入的完整链路。读完本文你将能基于sdk/packages/python/iii的公开 API 编写一个可连入 iii 引擎的 Python worker并用仓库源码与测试验证每一个方法的行为与边界条件。本文以 docs/0-17-0/sdk-reference/python-sdk.mdx及其渲染版本 python-sdk.mdx.skill.md为骨架结合 sdk/packages/python/iii/src/iii/iii.py 等实现文件、iii-example 示例与tests/测试用例深度展开。原文档注明其为“手写的公开面快照最终参考由 SDK 源码生成”因此本文以源码为准对齐细节。安装与包结构pip install iii-sdk安装后以iii为包名导入。SDK 依赖websockets12.0、pydantic2.0、opentelemetry-api1.25并捆绑同版本内部包iii-helpers见 pyproject.toml。Python 版本要求3.10。包根模块 iii/init.py 只导出公开面核心register_worker、TriggerAction、InitOptions、DEFAULT_ENGINE_URL、IIIConnectionState、ConnectionStateCallback、TelemetryOptions错误InvocationError、RegistrationRejectedError消息/队列TriggerActionEnqueue、TriggerActionVoid、EnqueueResult来自iii_helpers.queue类型与流IIIClient、StreamRequest、StreamResponse、IStream文档中提到但未从包根重新导出的类型需要从子模块显式导入符号所在模块说明WorkerInfoiii.iii_types存在但未从包根 re-exportIIIConnectionState字面量别名iii.iii_constants包根仅导出其类型别名用法RegisterFunctionFormatiii.iii_typesPython 专用格式助手TriggerHandler/TriggerConfigiii.trigger自定义触发器类型的抽象基类FunctionRefiii.iii_constants经iii.runtime汇总注册句柄ChannelReader/ChannelWriteriii.channels流式通道helpers自由函数iii.helperscreate_channel/create_stream等WorkerMetadata不属于本 SDK。FunctionInfo/TriggerInfo在 hooks.py 示例 中已不再作为 SDK 类存在引擎的engine::functions::list直接返回原始 dict 行。连接引擎register_workerdef register_worker(address: str | None None, options: InitOptions | None None) - III: ...address是引擎的 SDK WebSocket 地址如ws://localhost:49134。省略时按III_URL环境变量 →DEFAULT_ENGINE_URLws://127.0.0.1:49134见 iii_constants.py的优先级解析resolve_engine_url。返回的III实例承载下文所有方法。构造即自动在后台事件循环线程上发起connect_asyncregister_worker会阻塞等待最多 30 秒直到 WebSocket 建立_wait_until_connected超时后仅告警并返回实例连接在后台持续重试注册消息在连上后自动冲刷消息先进入上限 1000 条的本地队列见MAX_QUEUE_SIZE。连接状态可用add_connection_state_listener订阅回调立即以当前状态触发一次之后每次转移触发返回幂等的取消订阅函数。回调运行在 SDK 后台事件循环线程上保持轻量且不要在其中调用同步 SDK 方法会抛RuntimeError。InitOptions配置项InitOptions 以 dataclass 形式提供字段默认值说明worker_namehostname:pidIII_WORKER_NAME优先worker 显示名worker_descriptionNone面向人/LLM 的一行摘要出现在engine::workers::list/engine::workers::infonamespace无III_NAMESPACE兜底再缺省由引擎用defaultworker 及其函数注册、调用解析、触发器绑定的命名空间enable_metrics_reportingTrue通过 OpenTelemetry 上报 worker 指标invocation_timeout_ms30000trigger()默认超时reconnection_configReconnectionConfig()重连策略otel启用OtelConfig或 dict{enabled: False}或OTEL_ENABLEDfalse可关闭headersNone附加到 WebSocket 握手头的 dicttelemetryNoneTelemetryOptions上报 language/project/framework 等命名空间解析有一个值得注意的语义声明为空白字符串的 namespace 会抛ValueError而不是当作“未设置”_worker_namespace避免整个项目悄悄注册进错误的命名空间。重连配置 ReconnectionConfig 默认initial_delay_ms1000、max_delay_ms30000、backoff_multiplier2.0、jitter_factor0.3、max_retries-1无限重试。连接状态机与重连IIIConnectionState是字面量联合disconnected | connecting | connected | reconnecting | failed。重连循环实现指数退避并叠加随机抖动达到max_retries后进入failed_reconnect_loop。重连成功时SDK 会先发送reattach帧携带上次引擎分配的worker_id与reattach_token再重放全部触发器类型、函数与触发器注册最后冲刷排队消息_on_connected该顺序由测试 test_reconnect_sends_reattach.py 钉死避免重放与引擎清理旧连接发生竞态。注册拒绝与生命周期收尾引擎推送registrationrejected时FUNCTION_NAMESPACE_CONFLICT仅拒绝该函数worker 继续服务其余导出非致命记日志WORKER_NAMESPACE_CONFLICT或未知 code致命SDK 停止并不再重连抛出RegistrationRejectedError携带code/namespace/worker_name/owner_worker_id。收尾用shutdown()阻塞版或shutdown_async()asyncio 上下文取消重连与接收任务、以codeSHUTDOWN拒绝所有在途调用、关闭 WebSocket、停止后台线程并关闭 OTel。关闭后实例不可复用。注册函数register_functiondef 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: ...要点实现见 register_functionhandler 同步或异步均可同步 handler 会被包进 daemon 线程用run_in_executor等价方式执行避免阻塞事件循环异步 handler 直接 await。handler 签名第一个参数是触发 payloaddata若函数声明了字面名为metadata的参数如def handler(data, metadataNone)或def handler(data, *, metadataNone)SDK 会按位置/关键字传递每次调用的元数据未声明则保持handler(data)旧签名不变_metadata_passing_mode。格式自动提取Python 特色request_format/response_format缺省为None时从 handler 第一个参数与返回值的类型注解自动提取 JSON Schema显式传 schema 可覆盖无法注册“无 schema”的带类型 handler。Node SDK 因运行时擦除 TS 类型只能显式声明这是 Python SDK 的差异化能力。转换逻辑见 format_utils.py支持str/int/float/bool原语、Optional[X]生成[type,null]、list[X]、dict[str, X]、PydanticBaseModel复用model_json_schema()。HTTP 托管函数第二个参数也可传HttpInvocationConfig用于把函数指向 Lambda、Cloudflare Workers 等外部 HTTP 端点注册为 HTTP 调用的远程函数本地不可直接调用会回function_not_invokable。校验function_id非字符串抛TypeError空串或重复注册抛ValueError。返回值FunctionRefidunregister()调用unregister()发送unregisterfunction帧从引擎移除。from iii import register_worker, InitOptions from pydantic import BaseModel worker register_worker(ws://localhost:49134, InitOptions(worker_namemy-worker)) 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)触发器内置绑定与自定义类型绑定触发器register_triggerdef register_trigger(self, trigger: RegisterTriggerInput | dict[str, Any]) - Trigger: ...传入RegisterTriggerInput或等价 dict绑定一个触发器实例type使用的触发器类型 id如http、cron、state、subscribefunction_id触发时要调用的函数config类型相关配置HTTP 路径、cron 表达式等metadata每次触发都会传给 handler 的任意元数据namespace目标函数所在命名空间缺省取本 worker 的命名空间不是引擎default——因为函数注册在自己的命名空间里触发器默认落在default会永远解不到目标trigger_namespace触发器类型 provider 所在命名空间None表示引擎两步解析先本 worker 命名空间再引擎自身。触发器 id 由 SDK 自动生成 UUID。返回的Trigger句柄只有unregister()方法不存在顶层unregister_trigger。trigger worker.register_trigger({ type: http, function_id: greet, config: {api_path: /greet, http_method: GET}, }) trigger.unregister()自定义触发器类型register_trigger_type/unregister_trigger_typedef 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: ...RegisterTriggerTypeInput声明类型的id、description、trigger_request_format触发器配置的 JSON Schema / Pydantic 类与call_request_format触发时发给函数的 payload 格式。handler是 TriggerHandler 抽象基类实现register_trigger(config)/unregister_trigger(config)接收TriggerConfigid、function_id、config、metadata、namespace。返回的TriggerTypeRef是类型化句柄register_trigger(function_id, config)会用声明好的配置类校验并自动序列化且默认把触发器命名空间填为本 worker 命名空间register_function(function_id, handler)注册接收该类型调用载荷的函数。完整实例见 iii-example/src/trigger_types.pywebhook worker.register_trigger_type( RegisterTriggerTypeInput( idwebhook, descriptionIncoming webhook trigger, trigger_request_formatWebhookTriggerConfig, call_request_formatWebhookCallRequest, ), WebhookHandler(), ) async def handle_webhook(data: WebhookCallRequest) - dict: return {processed: True, method: data.method} webhook.register_function(example::webhook_handler, handle_webhook) webhook.register_trigger(example::webhook_handler, WebhookTriggerConfig(url/hooks/my-service))引擎侧registrationrejected推送会进入_handle_registration_rejectediii.py对函数级冲突仅告警跳过对 worker 级冲突置为致命错误。发起调用trigger/trigger_asyncdef trigger(self, request: dict[str, Any] | TriggerRequest) - Any: ... async def trigger_async(self, request: dict[str, Any] | TriggerRequest) - Any: ...trigger是同步入口把协程调度到 SDK 内部事件循环线程执行trigger_async供 asyncio 调用方 await。TriggerRequest字段function_id、payload、action、timeout_ms、metadata、namespace。返回值取决于actionaction语义返回值缺省同步请求-响应函数返回值TriggerAction.Enqueue(queuename)经命名队列异步路由需 worker-compose 中声明队列EnqueueResult含messageReceiptId的 dictTriggerAction.Void()fire-and-forgetNoneresult worker.trigger({function_id: greet, payload: {name: World}}) worker.trigger({function_id: notify, payload: {}, action: TriggerAction.Void()}) worker.trigger({function_id: process, payload: {}, action: TriggerAction.Enqueue(queuejobs)})实现细节同步无 action与 Enqueue 都会生成 UUIDinvocation_id并在self._pending登记 future用asyncio.wait_for限时等待超时抛codeTIMEOUT的InvocationError默认 30000ms可被请求级timeout_ms覆盖。Void 只发invokefunction帧不等待响应。每个请求自动注入traceparent与baggageW3C 上下文传播。命名空间解析显式namespace优先缺省时engine::前缀的内建函数走default其余继承本 worker 命名空间_invocation_namespace。TriggerAction.Enqueue的queue是仅限关键字的参数TriggerActionEnqueue(queueorders)序列化为{type:enqueue,queue:orders}TriggerActionVoid()序列化为{type:void}由测试 test_trigger_action.py 固定线格式。错误体系文档中的IIIInvocationError/IIIForbiddenError/IIITimeoutError在当前源码中已统一为单一基类InvocationErrorerrors.py旧名iii.IIIInvocationError不再存在于包根测试 test_errors.py 明确断言移除。分类依据是code字段而非子类场景InvocationError.code任意引擎侧调用错误引擎返回的 code如FORBIDDEN、TIMEOUT、HANDLER等RBAC 拒绝FORBIDDEN调用超时TIMEOUT畸形 wire 错误回退UNKNOWN基类属性code、message、function_id、stacktrace、invocation_id。str(err)只含code: message绝不包含 stacktrace远程引擎侧堆栈可能带内部路径不应暴露给最终用户_wrap_wire_error对任何畸形错误形状非 dict、缺字段、非字符串值都能安全兜底避免打印成裸 dict。InvocationError继承Exception现有except Exception捕获不受影响。另有两个专用异常RegistrationRejectedError注册被引擎致命拒绝见上文。流式 ChannelChannelReader/ChannelWriterChannelReader与ChannelWriter封装引擎的流式 WebSocket路径形如{base}/ws/channels/{id}?key{key}dirread|write见 channels.py 的build_channel_url。StreamChannelRef标识一个通道class StreamChannelRef(BaseModel): channel_id: str access_key: str direction: Literal[read, write]两个类都通过引擎 WS 基地址 StreamChannelRef构造ChannelWriter.write(bytes)二进制写单帧上限 64KB大载荷自动分帧close_async()关闭前先睡 10ms 让 TCP 栈冲刷缓冲避免关闭帧抢先于数据帧造成截断提供WritableStream兼容 Node 的.write/.end语义。ChannelReaderasync for chunk in reader迭代二进制块文本消息分发给on_message回调read_all()一次性读完整流。通道引用可以嵌套在 payload 中传递_resolve_channelsiii.py会在 handler 入参里把 dict 形态的StreamChannelRef递归替换成真实的ChannelReader/ChannelWriter。创建通道对用iii.helpers.create_channel([buffer_size])/create_channel_async(...)底层调用引擎内建engine::channels::create返回含 writer/reader/两个 ref 的Channel。集成测试 test_data_channels.py 展示了“发送方写通道、接收方read_all()计算统计”的完整 worker 间数据流test_api_triggers.py中还有 PDF 下载/上传与 SSE 流式用例。流Stream抽象IStreamstream.py声明get/set/delete/list/list_groups/update可用iii.helpers.create_stream(iii, name, impl)注册为stream::op(name)函数族update 由引擎内建原子更新逻辑处理不注册。日志与可观测性Logger实现于 iii_helpers/observability/logger.py经iii_helpers.observability暴露提供info/warn/error/debug每个方法接受消息 可选 data dictfrom iii_helpers.observability import Logger logger Logger(service_nameorders-api) logger.info(Order processed, {order_id: ord_123, status: completed}) logger.error(Payment failed, {order_id: ord_123, error_code: card_declined})日志以OpenTelemetry LogRecord形式发出自动携带当前 trace/span 上下文把日志与分布式链路关联起来OTel 未初始化时优雅回退到 Pythonlogging。结构化 data 会作为 OTel 属性log.data写入便于在 Grafana 等后端过滤聚合。遥测的初始化/导出侧属于iii-observabilityOtelConfig支持enabled、service_name、service_version、spans_flush_interval_ms默认 100ms比 OTel 默认 5s 更快出 trace、logs_enabled、metrics_enabled与engine_ws_url等telemetry_types.py。在 worker 侧每次 handler 调用会创建命名execute function_id的内部 span并记录iii.invocation.input/iii.invocation.output事件含 payload 红action 与截断受III_DISABLE_TRACE_PAYLOADS与字节上限环境变量控制见 iii.py 的_invoke_with_otel_context。消息帧与格式助手MessageTypeiii_types.py是运行时枚举命名 SDK 与引擎交换的每一种线帧registerfunction、unregisterfunction、invokefunction、invocationresult、registertriggertype、registertrigger、unregistertrigger、unregistertriggertype、triggerregistrationresult、workerregistered、registrationrejected、reattach等。它由中间件内部使用调用方很少需要直接接触。RegisterFunctionFormatiii_types.py是 Python 专用助手结构化声明函数请求/响应格式字段说明name参数名typestring/number/boolean/object/array/null/mapdescription参数说明bodyobject 类型的嵌套字段itemsarray 类型的元素 schemarequired是否必填直接传 JSON Schema dict 与构造RegisterFunctionFormat两种形式都汇入register_functionschema 随函数存储供 iii 控制台与 Agent 可读技能使用。信息类型与运维内建函数FunctionInfofunction_id、可选description、可选request_format/response_format、可选metadata。TriggerInfoid、trigger_type、function_id、可选config/metadata。WorkerInfo存在于iii.iii_types但未从包根 re-export需要时from iii.iii_types import WorkerInfoWorkerMetadata不属于本 SDK。引擎内建函数 id 常量见 iii_constants.py 的EngineFunctionsengine::functions::list、engine::functions::info、engine::workers::list、engine::workers::info、engine::triggers::list、engine::triggers::info、engine::registered-triggers::list、engine::registered-triggers::info、engine::workers::register与EngineTriggersengine::functions-available、log。iii-example 的 trigger_types.py 展示了用engine::triggers::list枚举全部可用触发器类型及其描述。完整可运行示例组合上文所有 API 的最小 worker风格参照 iii-example/src/main.pyimport os from iii import InitOptions, register_worker, TriggerAction from iii_helpers.observability import Logger logger Logger(service_namedemo-worker) engine_ws_url os.environ.get(III_URL, ws://localhost:49134) worker register_worker( addressengine_ws_url, optionsInitOptions( worker_namedemo-worker, otel{enabled: True, service_name: demo-worker}, ), ) def greet(data): return {message: fHello, {data[name]}!} worker.register_function(hello::greet, greet, descriptionGreets a user) worker.register_trigger({ type: http, function_id: hello::greet, config: {api_path: /greet, http_method: POST}, }) result worker.trigger({function_id: hello::greet, payload: {name: world}}) logger.info(greeted, {result: result}) # fire-and-forget worker.trigger({function_id: notify, payload: {}, action: TriggerAction.Void()}) worker.shutdown()要点回顾register_worker返回即建立或后台持续重试建立连接函数与触发器注册可先于连接完成SDK 会在连上后自动冲刷调用按action区分同步 / 入队 / void 三种语义错误统一以InvocationError.code分类处理流式数据通过iii.helpers.create_channel建立、ChannelReader/ChannelWriter读写日志与调用追踪天然接入 OpenTelemetry。若需深入了解某类方法的行为可继续阅读 sdk/packages/python/iii/src/iii/iii.py 与对应测试目录 sdk/packages/python/iii/tests 中的契约用例。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表