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/StreamResponse的http装饰器、以 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.http | HTTP 请求/响应类型、认证配置与http辅助函数 | http/init.py |
iii_helpers.observability | 结构化 Logger、OpenTelemetry 配置与 span 工具 | observability/init.py |
iii_helpers.queue | 队列入队结果类型 | queue/init.py |
iii_helpers.stream | Stream 触发器配置、变更事件、IO 输入输出与更新操作 | stream/init.py |
iii_helpers.worker_connection_manager | RBAC 认证与注册回调类型 | worker_connection_manager/init.py |
helpers/tests/下提供了test_observability_telemetry.py与test_observability_http_instrumentation.py两组测试,sdk/packages/python/iii/tests/下还有test_stream_models.py、test_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]])从源码看,http的wrapper会依次处理三种输入形态(见 http/init.py):
- 输入是
InternalHttpRequest实例:直接使用; - 输入是
dict:从path_params、query_params、body、headers、method、response、request_body键还原出InternalHttpRequest; - 其他类型:原样透传。
随后构造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)的调用参数:
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
url | str | 否 | 调用的目标 URL |
method | HttpMethod | 否 | HTTP 方法,默认'POST' |
timeout_ms | int \| None | 否 | 请求超时(毫秒) |
headers | dict[str, str] \| None | 否 | 附加请求头 |
auth | HttpAuthConfig \| None | 否 | 认证配置(bearer / HMAC / API key) |
其中HttpMethod是源码中定义的Literal["GET", "POST", "PUT", "PATCH", "DELETE"],注意它与核心builtin_triggers的 HTTP 方法枚举不同(后者还覆盖 HEAD/OPTIONS),见 http/init.py。
三种认证类型
认证字段的值都指向环境变量名,而不是明文密钥本身,敏感信息不会出现在代码或配置里:
| 类型 | 字段 | 说明 |
|---|---|---|
HttpAuthApiKey | header:API key 的自定义请求头名;value_key:存放 API key 值的环境变量名 | API key 通过自定义请求头发送 |
HttpAuthBearer | token_key:存放 bearer token 的环境变量名 | Bearer token 认证 |
HttpAuthHmac | secret_key:存放 HMAC 共享密钥的环境变量名 | 基于共享密钥的 HMAC 签名校验 |
三个模型都有type判别字段(Literal['api_key']/'bearer'/'hmac'),并由HttpAuthConfig = HttpAuthHmac | HttpAuthBearer | HttpAuthApiKey组成联合类型。
HttpRequest与HttpResponse
HttpRequest表示一个已缓冲的 HTTP 请求:
| 字段 | 类型 | 说明 |
|---|---|---|
method | str | 请求方法(如GET、POST) |
path_params | dict[str, str] | 从匹配路由提取的路径参数 |
query_params | dict[str, str \| list[str]] | URL 查询串参数 |
headers | dict[str, str \| list[str]] | 请求头 |
body | Any \| None | 已解析的请求体 |
HttpResponse表示要返回的缓冲响应:
| 字段 | 类型 | 说明 |
|---|---|---|
status_code | int | HTTP 状态码 |
headers | dict[str, str] | 响应头 |
body | Any \| None | 响应体 |
model_config | Any | Pydantic 模型配置 |
值得注意的实现细节: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):
| 字段 | 类型 | 默认值 / 说明 |
|---|---|---|
enabled | bool \| None | 是否启用 OTel,默认 True;设OTEL_ENABLED=false/0/no/off可关闭 |
service_name | str \| None | 服务名,默认取环境变量OTEL_SERVICE_NAME,否则'iii-python-sdk' |
service_version | str \| None | 服务版本,默认取SERVICE_VERSION,否则'unknown' |
service_namespace | str \| None | 服务命名空间属性 |
service_instance_id | str \| None | 服务实例 ID,默认随机 UUID |
engine_ws_url | str \| None | III 引擎 WebSocket 地址,默认取环境变量III_URL,否则'ws://localhost:49134' |
fetch_instrumentation_enabled | bool | 是否通过URLLibInstrumentor自动埋点 urllib HTTP 调用,默认 True |
spans_flush_interval_ms | int \| None | span 处理器刷盘延迟,默认 100ms;环境变量覆盖:OTEL_SPANS_FLUSH_INTERVAL_MS。文档特别提示:OpenTelemetry 默认的 5000ms 正是"操作结束后几秒才看到 trace"的原因 |
logs_enabled | bool \| None | 是否经EngineLogExporter导出 OTel 日志,OTel 启用时默认 True |
logs_flush_interval_ms | int \| None | 日志处理器刷盘延迟,默认 100ms |
logs_batch_size | int \| None | 每批导出的日志条数上限,默认 1 |
metrics_enabled | bool | 是否经EngineMetricsExporter导出指标,默认 True |
metrics_export_interval_ms | int | 指标导出间隔(毫秒),默认 60000(60 秒) |
ReconnectionConfig:WebSocket 重连策略
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
initial_delay_ms | int | 1000 | 起始延迟(毫秒) |
max_delay_ms | int | 30000 | 最大延迟上限(毫秒) |
backoff_multiplier | float | 2.0 | 指数退避乘数 |
jitter_factor | float | 0.3 | 随机抖动因子(0–1) |
max_retries | int | -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_ALLOWLIST | tuple[str, ...]类型的默认放行键集合(这些键不被脱敏) |
队列模块(iii_helpers.queue)
EnqueueResult是函数以TriggerAction.Enqueue方式被调用时返回的结果类型:
| 字段 | 类型 | 说明 |
|---|---|---|
messageReceiptId | str | 引擎为入队任务分配的 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_name | str | 是 | 要监听的流名称,只有该流上的变更触发 handler |
group_id | str \| None | 否 | 设置后仅该组内的变更触发 |
item_id | str \| None | 否 | 设置后仅该条目的变更触发 |
condition_function_id | str \| None | 否 | 条件函数 ID,返回 False 时跳过 handler |
StreamJoinLeaveTriggerConfig用于stream:join/stream:leave触发器,仅含condition_function_id一个可选字段。
Stream 事件类型
StreamChangeEvent是stream触发器的 handler 输入,在条目通过stream::set/stream::update/stream::delete变更时触发:
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
type | Literal['stream'] | 是 | 事件类型 |
timestamp | int | 是 | 事件的 Unix 时间戳 |
streamName | str | 是 | 变更发生的流 |
groupId | str | 是 | 变更发生的组 |
id | str \| None | 否 | 变更的条目 ID |
event | StreamChangeEventDetail | 是 | 变更详情 |
StreamChangeEventDetail包含type: Literal['create', 'update', 'delete'](变更类型)与data: Any(关联数据)。
StreamJoinLeaveEvent表示流的加入/离开事件,字段包括subscription_id(唯一订阅标识)、stream_name、group_id、id(可选条目 ID)与context(来自StreamAuthResult的认证上下文)。
认证相关类型:
| 类型 | 字段 | 说明 |
|---|---|---|
StreamAuthInput | headers: dict[str, str]、path: str、query_params: dict[str, list[str]]、addr: str(均必填) | 流认证输入 |
StreamAuthResult | context: Any \| None | 认证通过后传给 stream handler 的任意上下文 |
StreamJoinResult | unauthorized: bool(必填) | 加入是否未授权 |
Stream IO 输入输出
| 操作 | 输入字段 | 结果字段 |
|---|---|---|
| Get | stream_name、group_id、item_id | — |
| Set | stream_name、group_id、item_id、data: Any | old_value: TData \| None、new_value: TData |
| Delete | stream_name、group_id、item_id | old_value: Any \| None |
| List | stream_name、group_id | — |
| ListGroups | stream_name | — |
| Update | stream_name、group_id、item_id、ops: list[UpdateOp](原子应用的有序操作列表) | old_value、new_value、errors: list[UpdateOpError] |
更新操作(Update Ops)
UpdateOp是六种操作的联合类型:UpdateSet | UpdateIncrement | UpdateDecrement | UpdateAppend | UpdateRemove | UpdateMerge。
| 操作 | 类型标识 | 字段 | 语义 |
|---|---|---|---|
UpdateSet | "set" | path: str、value: Any | 把路径字段设为指定值;空字符串表示根值 |
UpdateIncrement | "increment" | path: str、by: int \| float | 数值字段增加by |
UpdateDecrement | "decrement" | path: str、by: int \| float | 数值字段减少by |
UpdateRemove | "remove" | path: str | 删除路径字段 |
UpdateAppend | "append" | path: MergePath \| None、value: Any | 数组追加 / 字符串拼接 / 嵌套路径 push |
UpdateMerge | "merge" | path: MergePath \| None、value: 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的顶层键覆盖同名键,兄弟键保留)。
安全与合法性校验(重点)
merge与append都会做严格的输入校验,违规操作返回结构化错误,且该操作不会生效:
- 路径深度 > 32 段;
- 段长度 > 256 字节;
- 值深度 > 16(
merge); - 顶层键 > 1024 个(
merge); - 任何
__proto__/constructor/prototype段或顶层键(防原型污染)。
错误通过state::update/stream::update响应的errors数组返回;成功应用的操作仍反映在响应的new_value中。
UpdateOpError
每个失败操作的错误详情:
| 字段 | 类型 | 说明 |
|---|---|---|
op_index | int | 出错操作在原始ops数组中的下标 |
code | str | 稳定错误码,如"merge.path.too_deep" |
message | str | 可读的错误描述(含具体数值) |
doc_url | str \| None | 可选的文档链接 |
实现细节:UpdateAppend与UpdateMerge都通过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 提供认证输入/输出与三类注册钩子的类型定义,是构建多租户、细粒度权限控制的基础。
AuthInput与AuthResult
AuthInput是 WebSocket 升级时传给 RBAC 认证函数的输入,包含升级请求的 HTTP 头、查询参数与客户端 IP:
| 字段 | 类型 | 说明 |
|---|---|---|
headers | dict[str, str] | WebSocket 升级请求的 HTTP 头 |
query_params | dict[str, list[str]] | 升级 URL 的查询参数,每个键对应值列表以支持重复键 |
ip_address | str | 连接客户端 IP |
AuthResult控制认证后的 Worker 可以调用哪些函数、注册哪些触发器,以及向中间件转发什么上下文:
| 字段 | 类型 | 说明 |
|---|---|---|
allowed_functions | list[str] | 在expose_functions之外额外允许的函数 ID |
forbidden_functions | list[str] | 即使匹配expose_functions也拒绝的函数 ID |
allowed_trigger_types | list[str] \| None | 允许注册触发器的触发器类型 ID,None表示全部允许 |
allow_trigger_type_registration | bool | 是否允许注册新的触发器类型 |
allow_function_registration | bool | 是否允许注册新函数 |
function_registration_prefix | str \| None | 应用于该 Worker 注册的所有函数 ID 的前缀 |
context | dict[str, Any] | 每次调用转发给中间件的任意上下文 |
从源码看(见 worker_connection_manager/init.py),AuthResult还支持namespaces: dict[str, list[str]]字段,用于按命名空间授予作用域化权限,例如{"orders": ["svc::*"]};值可以是精确函数 ID,也可以是通配符(支持裸写或match("...")写法)。当namespaces为空时,会话可声明任意命名空间,仅受allowed_functions与expose_functions约束。
注册钩子
三类钩子分别在 Worker 通过 RBAC 端口注册函数、触发器、触发器类型时被调用,返回映射后的结果,或抛异常拒绝注册。省略的字段保留注册请求中的原始值。
函数注册(on_function_registration_function_id):
OnFunctionRegistrationInput:function_id、description、metadata、context(会话的认证上下文);OnFunctionRegistrationResult:function_id、description、metadata(均为可选,省略即保留原值)。
触发器注册(on_trigger_registration_function_id):
OnTriggerRegistrationInput:trigger_id、trigger_type、function_id、config(触发器专属配置)、metadata、context;OnTriggerRegistrationResult:trigger_id、trigger_type、function_id、config(均为可选映射结果)。
触发器类型注册(on_trigger_type_registration_function_id):
OnTriggerTypeRegistrationInput:trigger_type_id、description、context;OnTriggerTypeRegistrationResult:trigger_type_id、description(均为可选映射结果)。
从源码还可看到,函数与触发器注册输入都包含namespace字段(显式注册值,缺省为"default");由于同一 ID 可以存在于多个命名空间,钩子需要依据命名空间逐项授权。
小结:三类落地场景
综合以上五个模块,iii-helpers在 III 项目中的典型组合方式是:
- HTTP 服务:用
http()包装 handler 获得强类型请求/响应,配合HttpInvocationConfig与三种HttpAuth*配置调用外部 HTTP 服务,密钥全部走环境变量; - 可观测性:用
OtelConfig初始化 OTel,Logger输出与 trace 关联的结构化日志,with_span/set_current_span_attribute/execute_traced_request串联分布式调用,flush_otel在短进程退出前保证数据送达; - 状态与实时流:用
StreamTriggerConfig订阅变更,通过StreamSetInput/StreamUpdateInput(内含merge、append等原子操作与原型污染防护)读写流数据,并用worker_connection_manager的类型构建 RBAC 认证与注册策略,把多租户权限落到函数、触发器与命名空间三个粒度上。
如需继续深入,可阅读同版本配套文档 docs/0-20-0/sdk-reference/ 下的 SDK 指南,以及sdk/packages/python/iii/tests/下的test_stream_models.py、test_stream_types.py、test_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),仅供参考