891 lines
39 KiB
Python
891 lines
39 KiB
Python
"""核心冻结类型(M1 设计 §2;最内层,禁止 import 任何实现)。
|
|
|
|
`LLMResponse` 前 11 个字段与三参考项目逐字保序——它们的测试按位置构造
|
|
fake,字段顺序即公共承诺;新增字段只增不删且必带默认值。
|
|
"""
|
|
|
|
import dataclasses
|
|
import json
|
|
import math
|
|
import re
|
|
import uuid
|
|
from collections.abc import Callable, Mapping
|
|
from dataclasses import dataclass, field
|
|
from enum import StrEnum
|
|
from types import MappingProxyType
|
|
from typing import Any, Literal
|
|
|
|
from loguru import logger
|
|
|
|
_MISSING_DONE_DOMAIN = frozenset({"retry", "salvage"})
|
|
|
|
_EFFORT_FALLBACK_DOMAIN = frozenset({"error", "nearest"})
|
|
"""`SourceConfig.effort_fallback` 的值域: 请求档打空时报错,还是映射到最近的档。"""
|
|
|
|
_PROTECTED_OVERLAY_KEYS: Mapping[str, str] = MappingProxyType(
|
|
{
|
|
"model": "会让遥测记录的 model 与实际请求分叉,成本按错单价换算",
|
|
"messages": "会同时破坏缓存 key 与遥测的 messages 口径",
|
|
"stream": "会绕过流式活性看门狗(TTFT/inter-token 超时全部失效)",
|
|
"stream_options": "会丢 usage 帧,导致成本遥测归零、TPM 闸按预扣量结算失准",
|
|
}
|
|
)
|
|
"""禁止出现在采样参数覆盖层里的键: 它们由治理层拥有,被覆盖即击穿治理。"""
|
|
|
|
USAGE_SOURCES = frozenset({"measured", "estimated", "unavailable"})
|
|
"""usage_source 值域;仅约束库内生产侧取值,不在 frozen dataclass 上做运行时校验。"""
|
|
|
|
_EST_TOKENS_QUOTA_DIVISOR = 60
|
|
"""未显式配置时的预扣量除数: 假定一次调用约占一秒钟的 TPM 配额份额。"""
|
|
|
|
_TENANT_ID_MAX_LEN = 128
|
|
_META_MAX_KEYS = 16
|
|
_META_VALUE_MAX_LEN = 256
|
|
_META_RESERVED_PREFIX = "pg_"
|
|
_META_KEY_RE = re.compile(r"[a-z0-9_.]{1,64}")
|
|
"""调用方维度的形态上限(issue #11 §4.2)。
|
|
|
|
数值取自同类系统的量级(Loki labels 15 / Salesforce 自定义索引 25 /
|
|
Sentry tag 200 字符),非本项目实测;键字符集照搬 OTel semconv。
|
|
`pg_` 前缀留给库将来的内建维度——同类做法见 LangSmith 的 `ls_`、
|
|
Traceloop 的 `traceloop.`;本版库自身不写入任何该前缀的键。"""
|
|
|
|
|
|
def validate_request_overlay(overlay: Mapping[str, Any], *, origin: str) -> dict[str, Any]:
|
|
"""校验采样参数覆盖层并返回浅拷贝;origin 用于把错误指回配置/调用点。
|
|
|
|
两类校验缺一不可(issue #4 设计决策 B):保护键会击穿治理;不可 JSON
|
|
序列化的值会在 `CacheMW` 的降级 try **之外**抛裸 `TypeError`——那条路径
|
|
不属错误四分类、`TelemetryMW` 也不捕,结果是一行遥测都没有就崩了。
|
|
两者都在进洋葱之前收口,故抛裸 `ValueError`(调用方编程错误,不可重试)。
|
|
"""
|
|
# Phase 1: 键形态——必须先于序列化试探,否则非 str 键会因 sort_keys 的
|
|
# 比较失败被误报成"值不可序列化",把人指向错误的方向
|
|
for key in overlay:
|
|
if not isinstance(key, str):
|
|
raise ValueError(f"{origin} 的键必须是 str: {key!r}(canonical JSON 要求)")
|
|
# Phase 2: 保护键
|
|
for key, reason in _PROTECTED_OVERLAY_KEYS.items():
|
|
if key in overlay:
|
|
raise ValueError(f"{origin} 不得覆盖 {key!r}: {reason}")
|
|
# Phase 3: 值可序列化(缓存 key 与遥测列都要 json.dumps)
|
|
try:
|
|
json.dumps(dict(overlay), sort_keys=True, ensure_ascii=False)
|
|
except (TypeError, ValueError) as exc:
|
|
raise ValueError(
|
|
f"{origin} 的值必须可 JSON 序列化(如 numpy 标量请先转 float/int): {exc}"
|
|
) from exc
|
|
return dict(overlay)
|
|
|
|
|
|
def validate_caller_dimensions(
|
|
tenant_id: str | None,
|
|
meta: Mapping[str, Any] | None,
|
|
*,
|
|
origin: str,
|
|
) -> tuple[str | None, dict[str, Any]]:
|
|
"""校验调用方自定义维度并返回浅拷贝;origin 用于把错误指回调用点(issue #11 §4.2)。
|
|
|
|
一切超限**报错而非静默丢弃**(P5): 同类系统里 Langfuse 对超长 value 直接
|
|
丢掉,那会让调用方以为记上了而实际没有。报错点必须在进洋葱之前——洋葱内
|
|
的失败都被遥测层降级成 warning,校验放那里等于没有校验。
|
|
"""
|
|
_validate_tenant_id(tenant_id, origin)
|
|
if meta is None:
|
|
return tenant_id, {}
|
|
# 键形态须先于值校验: 非 str 键若拖到值校验之后,会以更晦涩的形态报出来
|
|
_validate_meta_keys(meta, origin)
|
|
_validate_meta_values(meta, origin)
|
|
return tenant_id, dict(meta)
|
|
|
|
|
|
def _validate_tenant_id(tenant_id: str | None, origin: str) -> None:
|
|
"""租户标识形态;首尾空白**拒绝而非 strip**(见函数体注释)。"""
|
|
if tenant_id is None:
|
|
return
|
|
if not isinstance(tenant_id, str):
|
|
raise ValueError(f"{origin} 的 tenant_id 必须是 str: {tenant_id!r}")
|
|
# 悄悄 strip 会让 " t1" 变成 "t1": 二者在 RLS policy 的等值比较下是两个
|
|
# 不同租户,替调用方改写值等于把它的行藏进另一个租户,且不报错
|
|
if tenant_id != tenant_id.strip():
|
|
raise ValueError(
|
|
f"{origin} 的 tenant_id 不得含首尾空白: {tenant_id!r}"
|
|
"(RLS 等值比较下它与去空白版本是两个租户)"
|
|
)
|
|
if not tenant_id:
|
|
raise ValueError(f"{origin} 的 tenant_id 不得为空串(空串是未归属行的哨兵值)")
|
|
if len(tenant_id) > _TENANT_ID_MAX_LEN:
|
|
raise ValueError(f"{origin} 的 tenant_id 超长(上限 {_TENANT_ID_MAX_LEN}): {len(tenant_id)}")
|
|
|
|
|
|
def _validate_meta_keys(meta: Mapping[str, Any], origin: str) -> None:
|
|
"""键形态与数量;键集合被假定为低基数且稳定,故收紧到 OTel semconv 字符集。"""
|
|
# 数量闸先于逐键校验: 这道闸要防的正是"整个请求体被塞进 meta"的形态,
|
|
# 那时逐键正则会先跑上万次才报出真正的原因,拖慢的恰是出错路径
|
|
if len(meta) > _META_MAX_KEYS:
|
|
raise ValueError(f"{origin} 的 meta 键数超限(上限 {_META_MAX_KEYS}): {len(meta)}")
|
|
for key in meta:
|
|
if not isinstance(key, str):
|
|
raise ValueError(f"{origin} 的 meta 键必须是 str: {key!r}")
|
|
if key.startswith(_META_RESERVED_PREFIX):
|
|
raise ValueError(
|
|
f"{origin} 的 meta 键 {key!r} 使用了保留前缀 {_META_RESERVED_PREFIX!r}"
|
|
"(留给库将来的内建维度,避免与调用方的键撞名)"
|
|
)
|
|
if not _META_KEY_RE.fullmatch(key):
|
|
raise ValueError(
|
|
f"{origin} 的 meta 键 {key!r} 不合法: 只允许小写字母/数字/下划线/点,长度 1-64"
|
|
)
|
|
|
|
|
|
def _validate_meta_values(meta: Mapping[str, Any], origin: str) -> None:
|
|
"""值只收扁平标量;非有限 float 必须挡在这里。
|
|
|
|
`json.dumps` 会把 `nan`/`inf` 写成 `NaN`/`Infinity` 字面量——不是合法 JSON,
|
|
PG 的 JSONB 拒收。放行则写入失败会被遥测的降级 try 吞成 warning,即把调用方
|
|
的输入错误转成静默丢遥测(Codex 审查推翻了初稿"序列化不可达"的论断)。
|
|
"""
|
|
for key, value in meta.items():
|
|
if not isinstance(value, (str, int, float, bool)):
|
|
raise ValueError(
|
|
f"{origin} 的 meta 值必须是 str/int/float/bool: {key}={value!r}"
|
|
"(嵌套结构请调用方自行序列化)"
|
|
)
|
|
if isinstance(value, float) and not math.isfinite(value):
|
|
raise ValueError(f"{origin} 的 meta 值不得是 nan/inf: {key}={value!r}(非合法 JSON)")
|
|
if isinstance(value, str) and len(value) > _META_VALUE_MAX_LEN:
|
|
raise ValueError(
|
|
f"{origin} 的 meta 值超长(上限 {_META_VALUE_MAX_LEN}): {key} 长 {len(value)}"
|
|
)
|
|
|
|
|
|
def merge_sampling(extra_body: Mapping[str, Any], sampling: Mapping[str, Any]) -> dict[str, Any]:
|
|
"""合并配置级与调用级采样参数;调用级优先(issue #4 设计决策 A)。"""
|
|
return {**extra_body, **sampling}
|
|
|
|
|
|
def canonical_sampling_json(merged: Mapping[str, Any]) -> str | None:
|
|
"""缓存 key 与遥测 sampling 列共用的序列化口径;空 mapping → None。"""
|
|
if not merged:
|
|
return None
|
|
return json.dumps(dict(merged), sort_keys=True, ensure_ascii=False)
|
|
|
|
|
|
class Effort(StrEnum):
|
|
"""推理强度档位的封闭词汇(设计 §3.1)。
|
|
|
|
取值直接写进请求体(`reasoning_effort` 等键),**改名即改变发出去的字节**,
|
|
且会进缓存 key 与遥测落库,历史数据会断层。
|
|
|
|
八档而非六档: `none`(不推理)与 `auto`(推理,档位由模型自定)必须同时存在。
|
|
`auto` 不可省——newapi 上 26 个可调用模型里有 9 个是**纯开关型**(qwen 五个、
|
|
MiniMax-M3、glm-5/5.1/4.6v),它们能开推理却没有强度档可填;没有 `auto` 就只
|
|
能拿某个强度档冒充"开",而那正是本次要修的病根(旧 `thinking_on` 硬编码
|
|
`medium`,可 `medium` 在 GLM/kimi/deepseek 的档位表里根本不存在)。
|
|
|
|
词汇取四家参考实现共同收敛的一套(cherry-studio 的 canonical selection、
|
|
OpenRouter 的 `supported_efforts`、LiteLLM 的 `reasoning_effort_levels`、
|
|
new-api 的 `relayconvert/reasoning`),不自创。
|
|
|
|
**枚举定义在最内层而非决策层**: 它是 `SourceConfig`/`ChatRequest`/
|
|
`LLMResponse` 的字段类型,放进 `thinking.py` 会让 `types.py` 反向 import
|
|
决策模块(P7 依赖铁律),与 `ThinkingObservation` 同一理由。
|
|
"""
|
|
|
|
NONE = "none"
|
|
AUTO = "auto"
|
|
MINIMAL = "minimal"
|
|
LOW = "low"
|
|
MEDIUM = "medium"
|
|
HIGH = "high"
|
|
XHIGH = "xhigh"
|
|
MAX = "max"
|
|
|
|
|
|
EFFORT_ORDER: tuple[Effort, ...] = (
|
|
Effort.NONE,
|
|
Effort.MINIMAL,
|
|
Effort.LOW,
|
|
Effort.MEDIUM,
|
|
Effort.HIGH,
|
|
Effort.XHIGH,
|
|
Effort.MAX,
|
|
)
|
|
"""由弱到强的强度序;`AUTO` **不在其中**——它是"由模型自定",在强弱轴上没有位置。
|
|
|
|
供能力表求"最省的开启档"与 `nearest` 映射取最近档。公开(非 `_` 前缀)是因为
|
|
`thinking.py` 要跨模块消费它,跨模块引用私有名是坏味道。
|
|
"""
|
|
|
|
|
|
def coerce_effort(raw: Any, *, origin: str) -> Effort:
|
|
"""把外部传入的档位**归一**成 `Effort`;非法值报 `ValueError` 并列全八档。
|
|
|
|
存在的理由是"归一化点必须在入口":库内一律用 `is Effort.NONE` 做身份比较
|
|
(枚举成员唯一,`is` 比 `==` 更能表达"就是这一档"),而 `Effort` 是 `StrEnum`
|
|
——下游从 JSON/配置/命令行读出来的天然是裸字符串,`"none" is Effort.NONE`
|
|
恒为假。不在入口归一,身份比较就会在**错误路径上**误判(把一致的配置判成
|
|
矛盾),随后拼错误文案时再 `.value` 抛 `AttributeError`,连承诺的 `ValueError`
|
|
都拿不到(2026-09-05 独立验证实测)。
|
|
|
|
故裸字符串**接受并归一**而非拒收: 拒收会把 `.env` 之外的两条装配路(工厂 /
|
|
构造函数全量注入,CLAUDE.md §4.5)口径劈成两半,而 `.env` 那条早已是"解析即
|
|
归一"。`strip().lower()` 与 `config._to_effort` 同口径,理由同样是配置里的
|
|
行尾空格与大写写法是常态,而档位取值本身没有大小写语义。
|
|
|
|
`origin` 指回具体的配置项或调用点: 档位在源级、请求级两处都能配,只说
|
|
"非法档位"要人自己去找是哪一处填错了。传空串表示调用方自己会补上下文
|
|
(`config._cast` 的 `配置 X 解析失败` 已经说了是哪个 env 键)。
|
|
"""
|
|
if isinstance(raw, Effort):
|
|
return raw
|
|
prefix = f"{origin}: " if origin else ""
|
|
listed = ", ".join(e.value for e in Effort)
|
|
if isinstance(raw, str):
|
|
try:
|
|
return Effort(raw.strip().lower())
|
|
except ValueError:
|
|
# 不 `from exc`: 枚举原生的 "'lowest' is not a valid Effort" 只是同一
|
|
# 件事的英文复述,链上去反而把可操作的那句挤到后面
|
|
raise ValueError(f"{prefix}非法推理档位 {raw!r};允许: {listed}") from None
|
|
raise ValueError(
|
|
f"{prefix}推理档位必须是 Effort 或其字面量字符串,"
|
|
f"收到 {type(raw).__name__}: {raw!r};允许: {listed}"
|
|
)
|
|
|
|
|
|
class ThinkingObservation(StrEnum):
|
|
"""一次调用中"推理是否真的发生"的裁定结果(issue #16/#17)。
|
|
|
|
三态**不可折叠为布尔**: `UNKNOWN` 是"本次无任何信号,判不出来",与
|
|
`ABSENT`("上游明确上报了未推理")语义不同。把前者折叠进后者,正是
|
|
`reasoning_tokens=None` 制造的那个歧义——库据此静默宣称"没推理",而实际
|
|
可能推理了且已计费(MiniMax-M3 非流式实测: completion 53 vs 关闭档 3,
|
|
推理正文与 usage 明细双双不回传)。
|
|
|
|
裁定由 `thinking.observe_thinking` 做,本类只是取值域。**枚举定义在最内层
|
|
而非决策层**: 它是 `LLMResponse` 的字段类型,放进 `thinking.py` 会让
|
|
`types.py` 反向 import 决策模块(P7 依赖铁律)。
|
|
|
|
取值进遥测落库,改名即造成历史数据断层。
|
|
"""
|
|
|
|
OBSERVED = "observed"
|
|
ABSENT = "absent"
|
|
UNKNOWN = "unknown"
|
|
|
|
|
|
CallOperation = Literal["chat", "embed", "recognize_text", "parse_layout"]
|
|
"""遥测 `operation` 列的值域: **公开方法**四值,由调用点给定。
|
|
|
|
与 `PolyGatewayError.operation`(HTTP 子操作,如 `download_result`)是**两个语义**,
|
|
不做自动转换;链路上任何位置都不得读 `exc.operation` 来填本列(设计 §5 I1/I2)。"""
|
|
|
|
CALL_OPERATIONS: tuple[CallOperation, ...] = ("chat", "embed", "recognize_text", "parse_layout")
|
|
|
|
EventKind = Literal["attempt", "cache_hit", "terminal_failure"]
|
|
"""一行遥测描述的事件形态;旧行 NULL,不回填。
|
|
|
|
终态行与 attempt 行**不是重复事实**(前者描述逻辑终态,后者描述单次尝试),
|
|
故禁止按 `error IS NOT NULL` 跨两类直接计失败调用次数(设计 §6/§8)。"""
|
|
|
|
EVENT_KINDS: tuple[EventKind, ...] = ("attempt", "cache_hit", "terminal_failure")
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class CallStats:
|
|
"""一次**公开调用**(而非单次尝试)的统计快照(设计 §3)。
|
|
|
|
四种响应各平铺三字段会立刻漂移,故收敛成单一对象并由包根导出。
|
|
第三方合成响应的 `None` 表示**未知**,不得伪造 0。
|
|
"""
|
|
|
|
logical_call_id: str
|
|
"""每次公开调用一个 UUID;重试、结构化重问、embedding 分批共享同一个。
|
|
|
|
不占用既有 `parent_call_id`(后者是调用方的业务关联,语义不变)。"""
|
|
|
|
attempts: int
|
|
"""准入后实际调用 transport 端口的次数;含免预算 429 与端口本地拒绝。
|
|
|
|
**不是 HTTP 请求条数**: OCR layout 的 POST + ZIP GET 在同一次 transport
|
|
调用内,计 1 次。缓存命中与空输入是合法的零尝试。"""
|
|
|
|
total_latency_ms: int
|
|
"""从输入校验通过到返回/异常传播前的单调时钟快照。
|
|
|
|
含缓存 IO、退避等待、准入等待、重问、分批与内联记账。
|
|
"总耗时减最后一次尝试耗时"**不等于**纯等待(含其他本地工作)。"""
|
|
|
|
|
|
class _CallContext:
|
|
"""私有可变逻辑调用上下文: 只持计数、单调时钟与终态去重位,不做 I/O。
|
|
|
|
**每调用一个实例**的单任务对象: chat 重试、结构化重问、embedding 分批
|
|
都在同一任务内串行推进,故计数无需锁。**严禁提升为 client 实例属性**
|
|
——那会让同一 client 的并发调用互相串掉计数与逻辑 ID(库铁律"纯 asyncio 中立"、
|
|
VT `evolve_llm = llm` 教训的同一形态)。
|
|
"""
|
|
|
|
__slots__ = ("_attempts", "_now", "_started", "_terminal_claimed", "logical_call_id")
|
|
|
|
def __init__(self, *, now: Callable[[], float]) -> None:
|
|
self.logical_call_id = str(uuid.uuid4())
|
|
self._now = now
|
|
self._started = now()
|
|
self._attempts = 0
|
|
self._terminal_claimed = False
|
|
|
|
def register_attempt(self) -> None:
|
|
"""transport 调用**前**登记一次尝试(含免预算 429 与端口本地拒绝)。
|
|
|
|
登记点在调用前而非成功后: 否则失败与取消的尝试会从计数里消失,
|
|
而那正是诊断时最需要看见的那几次。
|
|
"""
|
|
self._attempts += 1
|
|
|
|
def snapshot(self) -> CallStats:
|
|
"""同步冻结当前快照;**绝不 await**,可多次调用。"""
|
|
return CallStats(
|
|
logical_call_id=self.logical_call_id,
|
|
attempts=self._attempts,
|
|
total_latency_ms=int((self._now() - self._started) * 1000),
|
|
)
|
|
|
|
def claim_terminal(self) -> bool:
|
|
"""首次 `True`、其后 `False`: 保证每逻辑调用至多写一条终态行。"""
|
|
if self._terminal_claimed:
|
|
return False
|
|
self._terminal_claimed = True
|
|
return True
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class LLMResponse:
|
|
"""一次治理调用的统一响应(与三项目超集兼容,ARCH §5.1)。"""
|
|
|
|
content: str
|
|
thinking: str
|
|
model: str
|
|
provider: str
|
|
prompt_tokens: int
|
|
completion_tokens: int
|
|
latency_ms: int
|
|
ttft_ms: float | None
|
|
max_inter_token_ms: float | None
|
|
cache_hit: bool
|
|
"""**PolyGateway 自身响应缓存**命中(未产生网关调用);与供应商侧 prompt
|
|
cache 无关,后者见 `cached_prompt_tokens`。"""
|
|
call_id: str
|
|
# —— 库新增(只增不删,必带默认值;迁移兼容硬约束)——
|
|
source_name: str = ""
|
|
cost: float | None = None
|
|
usage_source: str = "measured"
|
|
structured_data: Any | None = None
|
|
cached_prompt_tokens: int | None = None
|
|
"""供应商 prompt cache 命中的输入 token 数(issue #3);None = 该源未上报,
|
|
与"上报了但是 0"(真实零命中)区分——两者对下游的处置不同。"""
|
|
model_reported: str | None = None
|
|
"""API 响应体里的 model 字段;None = 未上报。与 `model`(配置别名)可能
|
|
分叉——供应商把别名指向新权重时,实验复现必须认这个串。"""
|
|
reasoning_tokens: int | None = None
|
|
"""推理消耗的输出 token 数(含在 `completion_tokens` 内,故不影响成本总额,
|
|
只补归因;issue #6)。
|
|
|
|
`None` = **本次调用**未上报,**不是**"该源不上报"——中转网关在上游不返回
|
|
usage 时会用本地 tokenizer 补算并整体替换 usage 对象,把
|
|
`completion_tokens_details` 一并吃掉(findings §4c 实测同一请求 10 轮呈
|
|
6:4 双峰)。实测三家供应商在未推理时都是整个 details 缺失、无人上报 `0`,
|
|
故下游判据须为 `in (None, 0)`,写 `== 0` 的条件永远不成立。
|
|
|
|
**该口径 2026-08-25 作废**(issue #16/#17): 供应商可能整体停报
|
|
`completion_tokens_details`(MiniMax 这一路实测已停),此时 `None` 只意味着
|
|
「没上报」而非「没推理」——同一次调用里库拿得到 185 字符推理正文。判「有没有
|
|
推理」一律改读 `thinking_observation`,上面那段只用于解读本版之前的历史数据。"""
|
|
|
|
thinking_observation: ThinkingObservation = ThinkingObservation.UNKNOWN
|
|
"""本次调用"推理是否真的发生"的三态裁定(issue #16/#17)。
|
|
|
|
`UNKNOWN` = **本次无任何信号,判不出来**,**不是**"没推理"——把两者折叠
|
|
是 `reasoning_tokens=None` 制造的老歧义。典型来源: 非流式路径下部分模型
|
|
推理已计费却既不回传正文也不回传 `completion_tokens_details`(MiniMax-M3
|
|
实测开启档 completion 53 vs 关闭档 3),该档即为 `UNKNOWN`。
|
|
要判"确实没推理"只认 `ABSENT`(上游明确上报 0)。"""
|
|
|
|
applied_effort: Effort | None = None
|
|
"""本次调用**真正发出去**的推理档位(issue #20);`None` = 调用方未表态。
|
|
|
|
与 `ChatRequest.reasoning_effort`(请求档)可能分叉: 源上配了
|
|
`EFFORT_FALLBACK=nearest` 时,请求 `medium` 而模型只有 low/high/max,实际发
|
|
出的是 `low`。遥测按本字段分组,记请求档会把整行挂在一个从未发出过的档下。
|
|
|
|
`None` 不是"没推理": 库不表态时也不推定模型自己的默认档——"没看见"不许说成
|
|
"发生了"(同 `thinking_observation` 的 `UNKNOWN` 一脉)。"""
|
|
|
|
call_stats: CallStats | None = None
|
|
"""本次**逻辑调用**的统计快照(1.3.5);`None` = 未知,不得读成 0。"""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class ChatRequest:
|
|
"""洋葱内部流转的不可变请求;中间件用 dataclasses.replace 派生,禁止原地修改。"""
|
|
|
|
messages: list[dict[str, Any]]
|
|
session_id: str | None = None
|
|
parent_call_id: str | None = None
|
|
cache_salt: str | None = None
|
|
cache_namespace: str | None = None
|
|
structured: Any | None = None
|
|
stream: bool = True
|
|
overlay: dict[str, Any] = field(default_factory=dict)
|
|
sampling: Mapping[str, Any] = field(default_factory=dict)
|
|
"""调用方采样意图的快照,库内中间件**永不修改**(issue #4 设计决策 A)。
|
|
|
|
与 `overlay` 分开是因为后者会被结构化中间件注入 `response_format`,在洋葱
|
|
不同深度取值不同;缓存 key 与三个遥测入口需要一个跨层恒定的读取点,否则
|
|
同一列在不同行口径分叉。"""
|
|
|
|
# —— 调用方自定义维度(issue #11;追加在末尾,不扰动既有字段的位置构造)——
|
|
tenant_id: str | None = None
|
|
"""调用方的租户标识,进遥测的 `tenant_id` 真实列(issue #11)。
|
|
|
|
独立成字段而非混进 `meta`,因为它是唯一享有真实列待遇的维度——可挂 RLS、
|
|
可进复合索引。混在 `meta` 里则调用方拼错(`tenantId`)不会报错,只会静默
|
|
降级成一个普通维度,正是本 issue 抱怨的失败形态。"""
|
|
|
|
meta: Mapping[str, Any] = field(default_factory=dict)
|
|
"""调用方自定义维度的只读快照,库不解释其含义,库内中间件**永不修改**。
|
|
|
|
**不进缓存 key**: 租户隔离已由 `cache_namespace` 负责并已进 key(ARCH §7.5),
|
|
再进一次既重复又会让存量缓存全量冷启动;且 `meta` 承载的是审计维度而非
|
|
语义维度,同 messages 同 namespace 下换个 batch_id 不应导致 miss。"""
|
|
|
|
# —— 请求级推理档位(issue #20;追加在末尾,不扰动既有字段的位置构造)——
|
|
reasoning_effort: Effort | None = None
|
|
"""本次调用要求的推理档位,压过源级默认(设计 §4.2 的最高优先级层)。
|
|
|
|
`None` 是**不表态**(随源级配置),与 `Effort.NONE`("要求不推理")严格区分:
|
|
把前者读成后者会让一次没写档位的调用悄悄关掉源上配好的推理。
|
|
|
|
独立成字段而非塞进 `overlay`: `overlay` 是采样参数的直通层,库不解释其内容,
|
|
而档位要经能力表校验、要进缓存 key、要落遥测——混进直通层等于放弃这三样,
|
|
正是 issue #20 里下游手写 `extra_body` 绕过全部治理的那条路。"""
|
|
|
|
# —— 库内部逻辑调用上下文(1.3.5;追加在末尾,不扰动既有字段的位置构造)——
|
|
call_context: _CallContext | None = field(default=None, compare=False, repr=False)
|
|
"""库内部逻辑调用上下文;`None` = 库内现场构造的请求,遥测 `logical_call_id` 落 NULL。
|
|
|
|
`compare=False, repr=False` 不是洁癖: 进 `compare` 会让两个内容相同的请求因
|
|
"不是同一次调用"而不相等,进 `repr` 则把库内部件泄进调用方的日志。
|
|
|
|
洋葱各层经 `dataclasses.replace` 派生请求时保留**同一引用**(不是拷贝),
|
|
重试/重问/分批才能共享同一个逻辑 ID 与计数。"""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class Usage:
|
|
"""token 用量;OCR 等无计费调用填 0。"""
|
|
|
|
prompt_tokens: int
|
|
completion_tokens: int
|
|
usage_source: str = "measured"
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class SourceStats:
|
|
"""限流后端回读的单源即时指标(CHS ports.py 同款)。"""
|
|
|
|
inflight: int
|
|
rpm_used: int
|
|
tpm_used: int
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class TelemetryStatus:
|
|
"""遥测后端的可写状态快照;degraded 期间下游可据此对账(issue #15)。
|
|
|
|
不叫 `health`: 库内 `health` 一律指**源的健康度**(`OcrTransport.check_health`
|
|
探活、`SourceSelector.health` 成功率 EWMA),而这里描述的是"这个 recorder
|
|
现在能不能写、为什么不能、丢了多少",是状态不是评分(设计 §3.3)。
|
|
|
|
时长一律给**相对秒数**而非绝对时间戳: 库内的时钟是 monotonic,把它的读数
|
|
交给下游会与 wall clock 混淆成两个不可比的时间轴。
|
|
"""
|
|
|
|
degraded: bool
|
|
fatal: bool
|
|
"""True = 本进程内不可恢复(仅 DSN 不可解析一类配置级失败),需改配置并重启。"""
|
|
reason: str | None
|
|
"""降级原因;未降级为 None。"""
|
|
degraded_for_s: float | None
|
|
"""已降级时长;未降级为 None。"""
|
|
dropped_rows: int
|
|
"""累计丢弃行数;**进程生命周期内单调不减**——恢复不等于没丢过。"""
|
|
retry_after_s: float | None
|
|
"""距下次重新准备的秒数;fatal 或未降级为 None,冷却已到期为 0.0。"""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class TransportResult:
|
|
"""transport 单次原始调用的产物;治理字段由 RetryMW 补齐为 LLMResponse。"""
|
|
|
|
content: str
|
|
thinking: str
|
|
prompt_tokens: int
|
|
completion_tokens: int
|
|
usage_source: str
|
|
ttft_ms: float | None
|
|
max_inter_token_ms: float | None
|
|
raw: dict[str, Any]
|
|
# —— 可观测字段(issue #3/#6;带默认值,非 OpenAI 兼容的 transport 可不填)——
|
|
cached_prompt_tokens: int | None = None
|
|
model_reported: str | None = None
|
|
reasoning_tokens: int | None = None
|
|
thinking_observation: ThinkingObservation = ThinkingObservation.UNKNOWN
|
|
"""本次调用"推理是否真的发生"的裁定(issue #16/#17),由 transport 组装时填。
|
|
|
|
默认 `UNKNOWN` 而非 `ABSENT`: 不做裁定的 transport(OCR/embedding 等)沉默
|
|
时,不该替上游做出"没推理"这个它从未做过的声明。"""
|
|
|
|
applied_effort: Effort | None = None
|
|
"""本次调用真正发出去的推理档位(issue #20),由做注入的 transport 填。
|
|
|
|
只有做了注入的那一层知道它: `nearest` 映射后请求档与实际档分叉(请求
|
|
`medium` → 实发 `low`),中间件事后再算一遍必然算成请求档。默认 `None` 是
|
|
"未表态/不注入推理参数"(OCR、embedding 等 transport 沉默即此)。"""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class SourceConfig:
|
|
"""单个模型源的完整配置(CHS config.py 超集;不变式在构造期报错)。
|
|
|
|
限额闸 0 表示不启用;`enable_thinking` 三态: None=不注入(模型默认)、
|
|
True=注入开启参数、False=注入关闭参数(统一 VT 与 CHS 相反的现状)。
|
|
|
|
2026-09-04 起 `enable_thinking` 降级为 `reasoning_effort` 的语法糖
|
|
(`True`→`AUTO`、`False`→`NONE`),保留不删是因为它已被三项目消费
|
|
(迁移兼容约束,ARCH §5.1)。两个字段说的是同一件事,故矛盾即报错。
|
|
"""
|
|
|
|
name: str
|
|
provider: str
|
|
base_url: str
|
|
api_key: str
|
|
model: str
|
|
timeout_s: float
|
|
max_concurrency: int = 0
|
|
rpm: int = 0
|
|
tpm: int = 0
|
|
est_tokens: int = 0
|
|
ttft_timeout_s: float | None = None
|
|
inter_token_timeout_s: float | None = None
|
|
enable_thinking: bool | None = None
|
|
missing_done: str = "retry"
|
|
trust_env: bool = True
|
|
extra_body: Mapping[str, Any] = field(default_factory=dict)
|
|
"""本源恒定的采样参数(如 `temperature=0`),并入请求体(issue #4)。
|
|
|
|
优先级低于调用级 overlay。注: 本字段令 `SourceConfig` 不再 hashable
|
|
(加任何 mapping 字段的固有代价,裸 dict 亦然),库内无以源作 key 的写法;
|
|
要可变副本用 `dict(source.extra_body)`,要改字段用 `dataclasses.replace`。"""
|
|
|
|
reasoning_effort: Effort | None = None
|
|
"""本源默认的推理档位;None = 不表态(与 `Effort.NONE`「要求不推理」不同)。
|
|
|
|
裸字符串(`"low"`、`" LOW "`)也收,构造期由 `coerce_effort` 归一成 `Effort`,
|
|
非法值当场 `ValueError` 并列出八档;**构造完成后本字段一定是 `Effort`**,库内
|
|
的 `is Effort.NONE` 身份比较依赖这条不变式。
|
|
|
|
**追加在末尾**是硬要求:三项目的测试按位置构造 fake,插在中间会静默错位
|
|
(本模块头部 docstring 的字段保序约定)。"""
|
|
|
|
effort_fallback: str = "error"
|
|
"""请求档打空时的处置: `error`(默认,报错)或 `nearest`(映射到最近的档)。
|
|
|
|
默认报错的理由是钱: 一次静默的 `medium → max` 在 GLM-5.3 上是数倍账单
|
|
(P5「严禁默认值掩盖错误」)。值域在此把关而非交给 `resolve_thinking`——
|
|
后者对未知值是 fail-closed(按 `error` 处理),不会替配置兜错,漏判的结果
|
|
就是 `EFFORT_FALLBAK` 这种拼写错误静默失效。"""
|
|
|
|
def __post_init__(self) -> None:
|
|
self._validate_identity()
|
|
self._validate_gates()
|
|
self._validate_watchdog()
|
|
self._validate_thinking()
|
|
self._freeze_extra_body()
|
|
|
|
def effective_est_tokens(self) -> int:
|
|
"""TPM 入场预扣量: 显式配置优先,否则按 tpm 派生(设计 §2.2)。"""
|
|
if self.est_tokens > 0:
|
|
return self.est_tokens
|
|
if self.tpm > 0:
|
|
return max(1, self.tpm // _EST_TOKENS_QUOTA_DIVISOR)
|
|
return 0
|
|
|
|
def _validate_identity(self) -> None:
|
|
for attr in ("name", "provider", "base_url", "api_key", "model"):
|
|
if not getattr(self, attr).strip():
|
|
raise ValueError(f"SourceConfig.{attr} 不能为空")
|
|
if self.missing_done not in _MISSING_DONE_DOMAIN:
|
|
raise ValueError(
|
|
f"missing_done 必须是 {sorted(_MISSING_DONE_DOMAIN)}: {self.missing_done!r}"
|
|
)
|
|
|
|
def _validate_gates(self) -> None:
|
|
if self.timeout_s <= 0:
|
|
raise ValueError("timeout_s 必须 > 0")
|
|
for attr in ("max_concurrency", "rpm", "tpm", "est_tokens"):
|
|
if getattr(self, attr) < 0:
|
|
raise ValueError(f"SourceConfig.{attr} 不能为负(0 表示不启用)")
|
|
# 注: 不再强制 `tpm > 0 ⇒ est_tokens > 0`——预扣量由 effective_est_tokens()
|
|
# 自 tpm 派生,运维只需填供应商配额页上抄得到的 tpm(设计 §3.2 #1)
|
|
|
|
def _validate_watchdog(self) -> None:
|
|
# CHS config.py:66-82: 流式看门狗成对配置且 0 < inter < ttft < timeout_s
|
|
if (self.ttft_timeout_s is None) != (self.inter_token_timeout_s is None):
|
|
raise ValueError("ttft_timeout_s 与 inter_token_timeout_s 必须同时设置或同时缺省")
|
|
if self.ttft_timeout_s is not None and not (
|
|
0 < self.inter_token_timeout_s < self.ttft_timeout_s < self.timeout_s
|
|
):
|
|
raise ValueError("看门狗不变式要求 0 < inter_token < ttft < timeout_s")
|
|
|
|
def _validate_thinking(self) -> None:
|
|
"""推理两键的**归一化**、值域与互不矛盾(issue #20 设计 §4.2)。
|
|
|
|
归一化必须先于下面的矛盾判定: 判据用的是 `is Effort.NONE`,而本类是公共
|
|
入口,`reasoning_effort="none"` 这种裸字符串写法(从 JSON/配置读出来的
|
|
常态)会让它误判成矛盾,再拼文案时 `.value` 直接 `AttributeError`。同一
|
|
理由也适用于下游读侧——归一化后库内一律是 `Effort`,`is` 比较才安全。
|
|
|
|
矛盾**报错而非「后者赢」**: `enable_thinking` 与 `reasoning_effort` 表达的是
|
|
同一件事,静默取其一等于替下游猜它到底想要哪个,而猜错的代价是账单——
|
|
猜成开启就是白花钱,猜成关闭就是拿到一个没推理过的答案。
|
|
|
|
判据是「二者是否都在说关闭」: `enable_thinking is False` 与
|
|
`reasoning_effort is NONE` 必须同真同假。`True` + 某个开启档(如 `low`)
|
|
不算矛盾,那只是把同一件事说了两遍,且后者更精确。
|
|
"""
|
|
if self.reasoning_effort is not None:
|
|
# frozen dataclass 改字段走 object.__setattr__(同款先例: _freeze_extra_body)
|
|
object.__setattr__(
|
|
self,
|
|
"reasoning_effort",
|
|
coerce_effort(
|
|
self.reasoning_effort, origin=f"SourceConfig({self.name}).reasoning_effort"
|
|
),
|
|
)
|
|
if isinstance(self.effort_fallback, str):
|
|
# 与相邻的 `REASONING_EFFORT` 同口径: `.env` 里的行尾空格与大写写法是
|
|
# 常态,而 `nearest`/`error` 本身没有大小写语义。归一化放在值域校验的
|
|
# 同一处(而不是 env 解析处),三条配置路一并覆盖
|
|
object.__setattr__(self, "effort_fallback", self.effort_fallback.strip().lower())
|
|
if self.effort_fallback not in _EFFORT_FALLBACK_DOMAIN:
|
|
raise ValueError(
|
|
f"SourceConfig.effort_fallback(EFFORT_FALLBACK)非法值 "
|
|
f"{self.effort_fallback!r};允许: {sorted(_EFFORT_FALLBACK_DOMAIN)}"
|
|
)
|
|
if self.enable_thinking is None or self.reasoning_effort is None:
|
|
return
|
|
if (self.enable_thinking is False) != (self.reasoning_effort is Effort.NONE):
|
|
raise ValueError(
|
|
f"源 {self.name!r} 的 enable_thinking={self.enable_thinking} 与 "
|
|
f"reasoning_effort={self.reasoning_effort.value!r} 相互矛盾: "
|
|
f"enable_thinking 已是 reasoning_effort 的语法糖"
|
|
f"(True={Effort.AUTO.value}、False={Effort.NONE.value})。"
|
|
f"请只保留其中一个,或让两者语义一致"
|
|
)
|
|
|
|
def _freeze_extra_body(self) -> None:
|
|
"""校验后转只读视图: 装配完成的源不应再被就地改采样参数(设计决策 E)。"""
|
|
validated = validate_request_overlay(
|
|
self.extra_body, origin=f"SourceConfig({self.name}).extra_body"
|
|
)
|
|
object.__setattr__(self, "extra_body", MappingProxyType(validated))
|
|
|
|
|
|
def strip_unsupported_extra_body(sources: list[SourceConfig], *, path: str) -> list[SourceConfig]:
|
|
"""剥离非 chat 路径不消费的 `extra_body` 并 warning(issue #4 决策 G)。
|
|
|
|
剥离是必需的而非顺手清理: embedding 的 payload 硬编码 `{model, input}`、
|
|
MonkeyOCR 只发 multipart 表单,两者都不会把 `extra_body` 发出去;但遥测的
|
|
`sampling` 列会并上 `source.extra_body`,不剥离就等于**记录一个从未发出的
|
|
参数**——那是数据造假,污染的恰是事后复现的唯一依据。
|
|
|
|
选择 warning 放行而非报错: 这两条路径本无采样语义,配错的后果远轻于 chat
|
|
路径,不值得让下游整个装配起不来(2026-07-31 人类拍板)。
|
|
"""
|
|
stripped = []
|
|
for source in sources:
|
|
if source.extra_body:
|
|
logger.warning(
|
|
"{} 路径暂不支持 extra_body,源 {} 的该配置已被忽略"
|
|
"(需要 dimensions 等参数请提 issue): {}",
|
|
path,
|
|
source.name,
|
|
dict(source.extra_body),
|
|
)
|
|
source = dataclasses.replace(source, extra_body={})
|
|
stripped.append(source)
|
|
return stripped
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class RetryPolicy:
|
|
"""重试策略;max_attempts = 总尝试次数(含首次,M1 设计 §2.3 统一语义)。"""
|
|
|
|
max_attempts: int
|
|
backoff_base_s: float
|
|
backoff_max_s: float
|
|
|
|
def __post_init__(self) -> None:
|
|
if self.max_attempts < 1:
|
|
raise ValueError("max_attempts 必须 ≥ 1(含首次尝试)")
|
|
if self.backoff_base_s <= 0 or self.backoff_max_s < self.backoff_base_s:
|
|
raise ValueError("退避参数要求 0 < backoff_base_s ≤ backoff_max_s")
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class BreakerConfig:
|
|
"""熔断配置;probe_ttl_s 是半开探针租约时长(持有者死亡后自动回收)。
|
|
|
|
M2.5 双通道: fail_threshold 是连续失败通道;min_calls/fail_rate/window_s
|
|
是失败率通道(窗口样本 ≥ min_calls 且失败率 ≥ fail_rate 即开路,429 不入);
|
|
开路时长按重开次数指数递增,封顶 max_cooldown_s(设计 2026-07-21-m25)。
|
|
"""
|
|
|
|
fail_threshold: int
|
|
cooldown_s: float
|
|
probe_ttl_s: float
|
|
min_calls: int = 10
|
|
fail_rate: float = 0.6
|
|
window_s: float = 60.0
|
|
max_cooldown_s: float = 300.0
|
|
|
|
def __post_init__(self) -> None:
|
|
if self.fail_threshold < 1 or self.cooldown_s <= 0 or self.probe_ttl_s <= 0:
|
|
raise ValueError("熔断配置要求 fail_threshold ≥ 1 且 cooldown_s/probe_ttl_s > 0")
|
|
if self.min_calls < 1 or not (0.0 < self.fail_rate <= 1.0) or self.window_s <= 0:
|
|
raise ValueError("失败率通道要求 min_calls ≥ 1、0 < fail_rate ≤ 1、window_s > 0")
|
|
if self.max_cooldown_s < self.cooldown_s:
|
|
raise ValueError("max_cooldown_s 不得小于 cooldown_s(退避封顶低于初值)")
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class BackpressurePolicy:
|
|
"""背压配置;M1 仅使用 poll_interval_s,stall 判定 M2 启用。"""
|
|
|
|
stall_window_s: float
|
|
poll_interval_s: float
|
|
|
|
def __post_init__(self) -> None:
|
|
if self.stall_window_s <= 0 or self.poll_interval_s <= 0:
|
|
raise ValueError("背压参数必须 > 0")
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class GlobalLimits:
|
|
"""scope 级全局限额;0 表示该闸不启用。"""
|
|
|
|
max_concurrency: int
|
|
rpm: int
|
|
tpm: int
|
|
|
|
def __post_init__(self) -> None:
|
|
if self.max_concurrency < 0 or self.rpm < 0 or self.tpm < 0:
|
|
raise ValueError("全局限额不能为负(0 表示不启用)")
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class OcrLayoutElement:
|
|
"""版面单元(M3 设计 §3.1): 来自 `_middle.json` para_blocks 的带类型块。
|
|
|
|
type 为开放字符串(实测 table/image/text,不枚举锁死——零业务假设);
|
|
bbox 为 OCR 原生页面坐标 (x1, y1, x2, y2),几何映射留业务侧(D9)。
|
|
"""
|
|
|
|
type: str
|
|
bbox: tuple[float, float, float, float]
|
|
page_index: int
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class OcrTextResult:
|
|
"""一次治理 OCR 文本转录的统一响应(/ocr/text;M3 设计 §3.1)。
|
|
|
|
text 空串 = 合法"无文字";行过滤/去重/拼帧留业务侧(VT 迁移 §3)。
|
|
"""
|
|
|
|
text: str
|
|
source_name: str
|
|
usage: Usage # OCR 无计费: Usage(0, 0);耗时由 latency_ms 承载
|
|
latency_ms: int
|
|
call_id: str
|
|
raw: dict[str, Any]
|
|
call_stats: CallStats | None = None
|
|
"""本次逻辑调用的统计快照(1.3.5);`None` = 未知。"""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class OcrLayoutResult:
|
|
"""一次治理版面解析的统一响应(/parse → ZIP;M3 设计 §3.1)。
|
|
|
|
elements 空 = 合法"无元素";CHS 首表 = 首个 type=="table" 元素。
|
|
page_sizes 按 page_index 索引。
|
|
"""
|
|
|
|
elements: list[OcrLayoutElement]
|
|
page_sizes: list[tuple[float, float]]
|
|
source_name: str
|
|
usage: Usage
|
|
latency_ms: int
|
|
call_id: str
|
|
raw: dict[str, Any]
|
|
call_stats: CallStats | None = None
|
|
"""本次逻辑调用的统计快照(1.3.5);`None` = 未知。"""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class OcrTextTransportResult:
|
|
"""transport 单次 /ocr/text 调用产物;治理字段由 OcrClient 补齐。"""
|
|
|
|
text: str
|
|
raw: dict[str, Any]
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class OcrLayoutTransportResult:
|
|
"""transport 单次 /parse 两段调用产物;治理字段由 OcrClient 补齐。"""
|
|
|
|
elements: list[OcrLayoutElement]
|
|
page_sizes: list[tuple[float, float]]
|
|
raw: dict[str, Any]
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class EmbeddingTransportResult:
|
|
"""一次原始 embedding 调用的解析结果(M2 设计 §7.2;transport → client)。"""
|
|
|
|
vectors: list[list[float]]
|
|
dim: int
|
|
prompt_tokens: int
|
|
usage_source: str # measured | estimated | unavailable
|
|
raw: dict[str, Any]
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class EmbeddingResponse:
|
|
"""一次治理 embedding 调用的统一响应(多批合并;与输入等长保序)。"""
|
|
|
|
vectors: list[list[float]]
|
|
dim: int
|
|
model: str
|
|
provider: str
|
|
prompt_tokens: int
|
|
usage_source: str
|
|
latency_ms: int
|
|
call_id: str
|
|
source_name: str
|
|
cost: float | None = None
|
|
call_stats: CallStats | None = None
|
|
"""本次逻辑调用(含全部分批)的统计快照(1.3.5);`None` = 未知。"""
|