news 2026/9/14 16:45:23

iii-helpers Python 包 API 完全指南:HTTP、OpenTelemetry 可观测性、Stream 与 RBAC 辅助类型

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
iii-helpers Python 包 API 完全指南:HTTP、OpenTelemetry 可观测性、Stream 与 RBAC 辅助类型

iii-helpers Python 包 API 完全指南:HTTP、OpenTelemetry 可观测性、Stream 与 RBAC 辅助类型

【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii

iii-helpers是 III 项目为 Python 开发者提供的官方辅助库,它把引擎侧的基础设施能力封装成一组可直接 import 的类型与函数:把引擎下发的原始 dict 请求包装成强类型StreamRequest/StreamResponsehttp装饰器、以 OpenTelemetry LogRecord 为出口的结构化Logger、覆盖 traceparent/baggage 的追踪上下文工具、Stream 触发器的全套输入输出与原子更新操作(set/merge/append等),以及 Worker 连接管理所需的 RBAC 认证与注册钩子类型。读完本文,你将能独立为 III Worker 编写带类型标注的 HTTP 处理器、接入链路追踪与结构化日志、操作 Stream 数据并实现基于 RBAC 的连接鉴权。

本文以 docs/0-20-0/api-reference/helpers-python.mdx 为骨架,并结合sdk/packages/python/helpers下的真实源码展开讲解。该文档由docs/next/scripts/generate-api-docs.mts自动生成,正文源文件是sdk/packages/python/helpers/src各模块的 docstring,因此源码注释与文档内容一一对应,可直接对照阅读。

安装与包结构

安装只需一条命令:

pip install iii-helpers

安装后可以从五个子模块导入所需内容,与源码目录sdk/packages/python/helpers/src/iii_helpers/一一对应:

子模块用途源码位置
iii_helpers.httpHTTP 请求/响应类型、认证配置与http辅助函数http/init.py
iii_helpers.observability结构化 Logger、OpenTelemetry 配置与 span 工具observability/init.py
iii_helpers.queue队列入队结果类型queue/init.py
iii_helpers.streamStream 触发器配置、变更事件、IO 输入输出与更新操作stream/init.py
iii_helpers.worker_connection_managerRBAC 认证与注册回调类型worker_connection_manager/init.py

helpers/tests/下提供了test_observability_telemetry.pytest_observability_http_instrumentation.py两组测试,sdk/packages/python/iii/tests/下还有test_stream_models.pytest_stream_types.py等测试,可作为理解各类型行为边界的参考。

HTTP 辅助模块(iii_helpers.http

HTTP 模块解决的是"引擎把原始请求数据交给 Python 处理器"的适配问题:函数内部无需手工解析 dict,而是拿到带类型的请求/响应对象。

http():把流式 handler 包装成引擎可直接调用的函数

http接收一个回调(req, res) -> HttpResponse | None,返回一个 III 引擎可直接调用的异步函数。包装器负责把引擎下发的原始 dict(或InternalHttpRequest)转换成回调期望的强类型StreamRequest/StreamResponse对。

签名

http(callback: Callable[Awaitable[HttpResponse[Any] | None]])

从源码看,httpwrapper会依次处理三种输入形态(见 http/init.py):

  1. 输入是InternalHttpRequest实例:直接使用;
  2. 输入是dict:从path_paramsquery_paramsbodyheadersmethodresponserequest_body键还原出InternalHttpRequest
  3. 其他类型:原样透传。

随后构造StreamResponse(internal.response)StreamRequest(...)并调用用户回调。这意味着你只需关注类型化的HttpRequest视图,底层的 WebSocket 响应通道由包装器代为维护。

用法示例

from iii_helpers.http import http, HttpResponse @app.function("http_handler") @http async def handler(req, res): # req: StreamRequest(含 method / headers / path_params / query_params / body) return HttpResponse(status_code=200, body={"ok": True})

HttpInvocationConfig:HTTP 外部函数调用配置

该配置用于描述对"外部 HTTP 服务"(如 Lambda、Cloudflare Workers)的调用参数:

字段类型必填说明
urlstr调用的目标 URL
methodHttpMethodHTTP 方法,默认'POST'
timeout_msint \| None请求超时(毫秒)
headersdict[str, str] \| None附加请求头
authHttpAuthConfig \| None认证配置(bearer / HMAC / API key)

其中HttpMethod是源码中定义的Literal["GET", "POST", "PUT", "PATCH", "DELETE"],注意它与核心builtin_triggers的 HTTP 方法枚举不同(后者还覆盖 HEAD/OPTIONS),见 http/init.py。

三种认证类型

认证字段的值都指向环境变量名,而不是明文密钥本身,敏感信息不会出现在代码或配置里:

类型字段说明
HttpAuthApiKeyheader:API key 的自定义请求头名;value_key:存放 API key 值的环境变量名API key 通过自定义请求头发送
HttpAuthBearertoken_key:存放 bearer token 的环境变量名Bearer token 认证
HttpAuthHmacsecret_key:存放 HMAC 共享密钥的环境变量名基于共享密钥的 HMAC 签名校验

三个模型都有type判别字段(Literal['api_key']/'bearer'/'hmac'),并由HttpAuthConfig = HttpAuthHmac | HttpAuthBearer | HttpAuthApiKey组成联合类型。

HttpRequestHttpResponse

HttpRequest表示一个已缓冲的 HTTP 请求:

字段类型说明
methodstr请求方法(如GETPOST
path_paramsdict[str, str]从匹配路由提取的路径参数
query_paramsdict[str, str \| list[str]]URL 查询串参数
headersdict[str, str \| list[str]]请求头
bodyAny \| None已解析的请求体

HttpResponse表示要返回的缓冲响应:

字段类型说明
status_codeintHTTP 状态码
headersdict[str, str]响应头
bodyAny \| None响应体
model_configAnyPydantic 模型配置

值得注意的实现细节:HttpResponse的 Pydantic 配置为ConfigDict(populate_by_name=True, arbitrary_types_allowed=True),且status_code在序列化时使用别名statusCode,与引擎侧和 Rust SDK 的线缆格式保持一致(见 http/init.py)。

可观测性模块(iii_helpers.observability

该模块是 Python Worker 接入可观测性的核心,涵盖结构化日志、OpenTelemetry 初始化与 span 操作。

Logger:以 OTel LogRecord 为出口的结构化日志器

Logger把每条日志作为 OpenTelemetry LogRecord 发出,每次调用都会自动捕获当前 trace 与 span 上下文,让日志与分布式链路天然关联,无需手工拼接 trace id。当 OTel 未初始化时,Logger会优雅地回退到 Python 标准logging

源码实现(见 observability/logger.py)展示了两个关键设计:

  • 四级方法info/warn/error/debug映射到 OTelSeverityNumber:INFO=9、WARN=13、ERROR=17、DEBUG=5;
  • 日志属性中自动携带service.name,结构化数据以log.data属性整体附加;
  • trace_id / span_id 优先取显式构造参数,否则取当前 span 上下文。

推荐用法:把结构化数据作为第二个参数传入

from iii import Logger logger = Logger() # 基础日志:trace 上下文自动注入 logger.info('Worker connected') # 结构化上下文:便于在可观测后端过滤、聚合、建仪表盘 logger.info('Order processed', {'order_id': 'ord_123', 'amount': 49.99, 'currency': 'USD'}) logger.warn('Retry attempt', {'attempt': 3, 'max_retries': 5, 'endpoint': '/api/charge'}) logger.error('Payment failed', { 'order_id': 'ord_123', 'gateway': 'stripe', 'error_code': 'card_declined', })

官方文档特别强调:使用 dict 键值对而不是字符串插值,才能在 Grafana、Datadog 等后端做过滤与聚合。另外Logger构造函数还接受可选的trace_id/span_id/service_name,用于在缺少活动 span 时显式绑定上下文。

OtelConfig:OpenTelemetry 初始化配置

init_otel(config: OtelConfig | None = None, loop: None = None)负责初始化 OpenTelemetry,后续重复调用为 no-op。其配置字段(源码见 observability/telemetry_types.py):

字段类型默认值 / 说明
enabledbool \| None是否启用 OTel,默认 True;设OTEL_ENABLED=false/0/no/off可关闭
service_namestr \| None服务名,默认取环境变量OTEL_SERVICE_NAME,否则'iii-python-sdk'
service_versionstr \| None服务版本,默认取SERVICE_VERSION,否则'unknown'
service_namespacestr \| None服务命名空间属性
service_instance_idstr \| None服务实例 ID,默认随机 UUID
engine_ws_urlstr \| NoneIII 引擎 WebSocket 地址,默认取环境变量III_URL,否则'ws://localhost:49134'
fetch_instrumentation_enabledbool是否通过URLLibInstrumentor自动埋点 urllib HTTP 调用,默认 True
spans_flush_interval_msint \| Nonespan 处理器刷盘延迟,默认 100ms;环境变量覆盖:OTEL_SPANS_FLUSH_INTERVAL_MS。文档特别提示:OpenTelemetry 默认的 5000ms 正是"操作结束后几秒才看到 trace"的原因
logs_enabledbool \| None是否经EngineLogExporter导出 OTel 日志,OTel 启用时默认 True
logs_flush_interval_msint \| None日志处理器刷盘延迟,默认 100ms
logs_batch_sizeint \| None每批导出的日志条数上限,默认 1
metrics_enabledbool是否经EngineMetricsExporter导出指标,默认 True
metrics_export_interval_msint指标导出间隔(毫秒),默认 60000(60 秒)

ReconnectionConfig:WebSocket 重连策略

字段类型默认值说明
initial_delay_msint1000起始延迟(毫秒)
max_delay_msint30000最大延迟上限(毫秒)
backoff_multiplierfloat2.0指数退避乘数
jitter_factorfloat0.3随机抖动因子(0–1)
max_retriesint-1最大重试次数,-1表示无限重试

Span 工具函数

函数签名要点行为
current_trace_id()无参返回当前活动 trace_id 的 32 位十六进制字符串,不可用时返回 None
current_span_id()无参返回当前活动 span_id 的 16 位十六进制字符串,不可用时返回 None
current_span_is_recording()无参无活动 span 或采样器丢弃 span 时返回 False
with_span(name, fn, kind=None, traceparent=None)异步启动新 span 并在其中运行fn(span);tracer 未初始化时用 no-op span 调用fn,静默忽略属性/事件调用
set_current_span_attribute(key, value)同步给当前 span 设置属性;span 未 recording 时为 no-op
record_span_event(name, attrs=None)同步记录 span 事件;span 未 recording 时为 no-op
set_current_span_error(message)同步标记当前 span 错误;无活动 span 时为 no-op

追踪上下文注入与提取

跨服务传播追踪信息时使用 W3C 标准头:

函数行为
inject_traceparent()把当前 trace 上下文注入 W3Ctraceparent头字符串
extract_traceparent(traceparent: str)从 W3Ctraceparent头字符串提取 trace 上下文
inject_baggage()把当前 baggage 注入 W3Cbaggage头字符串
extract_baggage(baggage: str)从 W3Cbaggage头字符串提取 baggage

配套的BaggageSpanProcessor负责把 baggage 中的键值同步到 span 上,保证跨服务调用时业务上下文不丢失。

生命周期与带追踪的 HTTP 调用

  • init_otel(config=None, loop=None):初始化 OpenTelemetry,重复调用为 no-op;
  • shutdown_otel():同步关闭 OTel(尽力而为,不等待 WS 刷盘);
  • flush_otel()shutdown_otel的对立面——在不拆除提供者的前提下强制刷盘所有 OTel 提供者。适合在短生命周期进程退出前使用:既希望挂起的 spans/metrics/logs 被送达,又计划继续使用 OTel;
  • execute_traced_request(client, request):在 OTel CLIENT span 内执行 httpx 请求。具体行为包括:向出站请求头注入 W3Ctraceparent;在 span 上记录 HTTP 语义约定属性;对 status >= 400 的响应设置 ERROR span 状态;对网络级错误记录异常。

日志脱敏与截断

可观测性数据经常包含敏感字段,模块提供了脱敏工具:

函数说明
redact(value: Any)对值进行脱敏处理
redact_and_truncate(value: Any, max_bytes: Optional[int] = None)脱敏并按字节上限截断
resolve_max_bytes_from_env()从环境变量解析截断上限
DEFAULT_ALLOWLISTtuple[str, ...]类型的默认放行键集合(这些键不被脱敏)

队列模块(iii_helpers.queue

EnqueueResult是函数以TriggerAction.Enqueue方式被调用时返回的结果类型:

字段类型说明
messageReceiptIdstr引擎为入队任务分配的 UUID 回执 ID
from iii_helpers.queue import EnqueueResult # 在函数内把消息投递到队列后返回回执 return EnqueueResult(messageReceiptId=receipt_id)

Stream 模块(iii_helpers.stream

Stream 是 III 中"流式数据 + 状态"的核心抽象,该模块提供四类类型:触发器配置、变更事件、IO 输入输出、更新操作。

Stream 触发器配置

StreamTriggerConfig用于配置stream触发器,决定哪些条目变更会触发 handler:

字段类型必填说明
stream_namestr要监听的流名称,只有该流上的变更触发 handler
group_idstr \| None设置后仅该组内的变更触发
item_idstr \| None设置后仅该条目的变更触发
condition_function_idstr \| None条件函数 ID,返回 False 时跳过 handler

StreamJoinLeaveTriggerConfig用于stream:join/stream:leave触发器,仅含condition_function_id一个可选字段。

Stream 事件类型

StreamChangeEventstream触发器的 handler 输入,在条目通过stream::set/stream::update/stream::delete变更时触发:

字段类型必填说明
typeLiteral['stream']事件类型
timestampint事件的 Unix 时间戳
streamNamestr变更发生的流
groupIdstr变更发生的组
idstr \| None变更的条目 ID
eventStreamChangeEventDetail变更详情

StreamChangeEventDetail包含type: Literal['create', 'update', 'delete'](变更类型)与data: Any(关联数据)。

StreamJoinLeaveEvent表示流的加入/离开事件,字段包括subscription_id(唯一订阅标识)、stream_namegroup_idid(可选条目 ID)与context(来自StreamAuthResult的认证上下文)。

认证相关类型:

类型字段说明
StreamAuthInputheaders: dict[str, str]path: strquery_params: dict[str, list[str]]addr: str(均必填)流认证输入
StreamAuthResultcontext: Any \| None认证通过后传给 stream handler 的任意上下文
StreamJoinResultunauthorized: bool(必填)加入是否未授权

Stream IO 输入输出

操作输入字段结果字段
Getstream_namegroup_iditem_id
Setstream_namegroup_iditem_iddata: Anyold_value: TData \| Nonenew_value: TData
Deletestream_namegroup_iditem_idold_value: Any \| None
Liststream_namegroup_id
ListGroupsstream_name
Updatestream_namegroup_iditem_idops: list[UpdateOp](原子应用的有序操作列表)old_valuenew_valueerrors: list[UpdateOpError]

更新操作(Update Ops)

UpdateOp是六种操作的联合类型:UpdateSet | UpdateIncrement | UpdateDecrement | UpdateAppend | UpdateRemove | UpdateMerge

操作类型标识字段语义
UpdateSet"set"path: strvalue: Any把路径字段设为指定值;空字符串表示根值
UpdateIncrement"increment"path: strby: int \| float数值字段增加by
UpdateDecrement"decrement"path: strby: int \| float数值字段减少by
UpdateRemove"remove"path: str删除路径字段
UpdateAppend"append"path: MergePath \| Nonevalue: Any数组追加 / 字符串拼接 / 嵌套路径 push
UpdateMerge"merge"path: MergePath \| Nonevalue: Any把对象浅合并到目标节点

MergePath = str | list[str],即"单字符串(旧式/一级键)或字面量段列表(嵌套路径)",与 Node SDK 的string | string[]别名及 Rust 的MergePath枚举对齐。

UpdateAppend的路径形式与引擎语义

接受的路径形式(与UpdateMerge一致):

  • None/""/[]:在根节点追加;
  • "foo":在一级键foo处追加。注意带点的字符串如"a.b"字面量键名"a.b",不会被拆解成a -> b
  • ["a", "b", "c"]:嵌套路径,每个元素都是字面量段。

引擎在叶子节点的行为:

  • 嵌套路径上的缺失/非对象中间节点自动以{}创建/替换;
  • 叶子缺失/null + 嵌套路径 →[value](总是数组);
  • 叶子缺失/null + 单字符串路径 → 字符串拼接层级按字符串处理,否则[value]
  • 已存在数组 → push;
  • 已存在字符串 + 字符串值 → 拼接;
  • 叶子是对象/标量 → 返回append.type_mismatch
UpdateMerge的路径形式与引擎语义
  • 根合并:None/""/[]
  • "foo"等价于["foo"],即一级键;
  • ["a", "b", "c"]:嵌套路径,每个元素是字面量键;["a.b"]写入的是名为"a.b"的单个键,而不是a -> b

引擎行为:路径上的缺失/非对象中间节点自动替换为{};合并是目标节点处的浅合并(value的顶层键覆盖同名键,兄弟键保留)。

安全与合法性校验(重点)

mergeappend都会做严格的输入校验,违规操作返回结构化错误,且该操作不会生效

  • 路径深度 > 32 段;
  • 段长度 > 256 字节;
  • 值深度 > 16(merge);
  • 顶层键 > 1024 个(merge);
  • 任何__proto__/constructor/prototype段或顶层键(防原型污染)。

错误通过state::update/stream::update响应的errors数组返回;成功应用的操作仍反映在响应的new_value中。

UpdateOpError

每个失败操作的错误详情:

字段类型说明
op_indexint出错操作在原始ops数组中的下标
codestr稳定错误码,如"merge.path.too_deep"
messagestr可读的错误描述(含具体数值)
doc_urlstr \| None可选的文档链接

实现细节:UpdateAppendUpdateMerge都通过model_serializer(mode="wrap")在序列化时省略path: None,从而让线缆载荷与 Rust SDK 的#[serde(skip_serializing_if = "Option::is_none")]逐字节一致(见 stream/init.py)。StreamUpdateResult.errors为空时也会从 JSON 线缆中省略该字段。

Worker 连接管理模块(RBAC)

该模块为通过 RBAC 端口连接的 Worker 提供认证输入/输出与三类注册钩子的类型定义,是构建多租户、细粒度权限控制的基础。

AuthInputAuthResult

AuthInput是 WebSocket 升级时传给 RBAC 认证函数的输入,包含升级请求的 HTTP 头、查询参数与客户端 IP:

字段类型说明
headersdict[str, str]WebSocket 升级请求的 HTTP 头
query_paramsdict[str, list[str]]升级 URL 的查询参数,每个键对应值列表以支持重复键
ip_addressstr连接客户端 IP

AuthResult控制认证后的 Worker 可以调用哪些函数、注册哪些触发器,以及向中间件转发什么上下文:

字段类型说明
allowed_functionslist[str]expose_functions之外额外允许的函数 ID
forbidden_functionslist[str]即使匹配expose_functions也拒绝的函数 ID
allowed_trigger_typeslist[str] \| None允许注册触发器的触发器类型 ID,None表示全部允许
allow_trigger_type_registrationbool是否允许注册新的触发器类型
allow_function_registrationbool是否允许注册新函数
function_registration_prefixstr \| None应用于该 Worker 注册的所有函数 ID 的前缀
contextdict[str, Any]每次调用转发给中间件的任意上下文

从源码看(见 worker_connection_manager/init.py),AuthResult还支持namespaces: dict[str, list[str]]字段,用于按命名空间授予作用域化权限,例如{"orders": ["svc::*"]};值可以是精确函数 ID,也可以是通配符(支持裸写或match("...")写法)。当namespaces为空时,会话可声明任意命名空间,仅受allowed_functionsexpose_functions约束。

注册钩子

三类钩子分别在 Worker 通过 RBAC 端口注册函数、触发器、触发器类型时被调用,返回映射后的结果,或抛异常拒绝注册。省略的字段保留注册请求中的原始值

函数注册on_function_registration_function_id):

  • OnFunctionRegistrationInputfunction_iddescriptionmetadatacontext(会话的认证上下文);
  • OnFunctionRegistrationResultfunction_iddescriptionmetadata(均为可选,省略即保留原值)。

触发器注册on_trigger_registration_function_id):

  • OnTriggerRegistrationInputtrigger_idtrigger_typefunction_idconfig(触发器专属配置)、metadatacontext
  • OnTriggerRegistrationResulttrigger_idtrigger_typefunction_idconfig(均为可选映射结果)。

触发器类型注册on_trigger_type_registration_function_id):

  • OnTriggerTypeRegistrationInputtrigger_type_iddescriptioncontext
  • OnTriggerTypeRegistrationResulttrigger_type_iddescription(均为可选映射结果)。

从源码还可看到,函数与触发器注册输入都包含namespace字段(显式注册值,缺省为"default");由于同一 ID 可以存在于多个命名空间,钩子需要依据命名空间逐项授权。

小结:三类落地场景

综合以上五个模块,iii-helpers在 III 项目中的典型组合方式是:

  1. HTTP 服务:用http()包装 handler 获得强类型请求/响应,配合HttpInvocationConfig与三种HttpAuth*配置调用外部 HTTP 服务,密钥全部走环境变量;
  2. 可观测性:用OtelConfig初始化 OTel,Logger输出与 trace 关联的结构化日志,with_span/set_current_span_attribute/execute_traced_request串联分布式调用,flush_otel在短进程退出前保证数据送达;
  3. 状态与实时流:用StreamTriggerConfig订阅变更,通过StreamSetInput/StreamUpdateInput(内含mergeappend等原子操作与原型污染防护)读写流数据,并用worker_connection_manager的类型构建 RBAC 认证与注册策略,把多租户权限落到函数、触发器与命名空间三个粒度上。

如需继续深入,可阅读同版本配套文档 docs/0-20-0/sdk-reference/ 下的 SDK 指南,以及sdk/packages/python/iii/tests/下的test_stream_models.pytest_stream_types.pytest_rbac_workers.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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/14 16:45:13

企业级Agent平台落地指南:从超级个体到超级团队

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/14 16:38:49

工业监控上位机开发:MVP架构与Modbus TCP优化实践

1. 工业监控上位机开发背景与需求 在自动化生产线中,力位移曲线监控是质量检测的核心环节。以汽车零部件压装工艺为例,每秒钟需要采集上百个压力传感器和位移传感器的数据点,通过实时曲线比对确保装配质量。传统方式依赖专用仪器,…

作者头像 李华
网站建设 2026/9/14 16:38:35

Flutter与鸿蒙结合:pro_mpack二进制序列化优化实践

1. 项目背景与核心价值 在鸿蒙生态快速发展的当下,Flutter作为跨平台开发框架与鸿蒙系统的结合越来越紧密。pro_mpack作为Flutter生态中的高性能二进制序列化库,其鸿蒙化适配对于提升分布式场景下的数据传输效率具有关键意义。MessagePack协议相比传统JS…

作者头像 李华