Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 166b2865d0 | |||
| b0ab39edbf | |||
| a07a6b096b | |||
| 53d2f089e1 | |||
| 5e2c35a812 | |||
| b1bc06e2b1 | |||
| da77b123ec | |||
| 9474c76ab0 | |||
| 1ff83bbfe0 | |||
| 6c640fcca3 |
@@ -69,6 +69,11 @@ LLM_CIRCUIT_BREAKER_COOLDOWN=60 # 或 LLM__BREAKER__COOLDOWN_S
|
|||||||
# ── 在几毫秒内死掉且 MAX_ATTEMPTS 一格用不上。wait 不削弱保护(等待期照样
|
# ── 在几毫秒内死掉且 MAX_ATTEMPTS 一格用不上。wait 不削弱保护(等待期照样
|
||||||
# ── 不发请求),只是把最坏墙钟拉长到 BACKPRESSURE__STALL_WINDOW_S ──
|
# ── 不发请求),只是把最坏墙钟拉长到 BACKPRESSURE__STALL_WINDOW_S ──
|
||||||
# LLM__CIRCUIT_OPEN=fail_fast # 熔断开路: fail_fast(默认) | wait
|
# LLM__CIRCUIT_OPEN=fail_fast # 熔断开路: fail_fast(默认) | wait
|
||||||
|
# LLM__CALL_DEADLINE_S= # 一次逻辑调用的墙钟硬边界(秒);缺省不设 = 不启用
|
||||||
|
# ── 治理对象是"等待"(退避/配额轮询/熔断冷却/结构化重问/embedding 分批共享一份),
|
||||||
|
# ── 不是单次 HTTP 超时(那是 TIMEOUT_S)。清理仍在 finally 跑完: 返回时刻 = 期限 + 清理耗时,
|
||||||
|
# ── 且到期 ≠ 未产出、≠ 未计费。到期抛 CallDeadlineExceeded(不属四分类、
|
||||||
|
# ── 不属 GatewayUnavailableError 族、无 retry_after_s);非法值(0/负/nan/inf)装配期报错 ──
|
||||||
|
|
||||||
# ══ 装配选择(PGW_*)══
|
# ══ 装配选择(PGW_*)══
|
||||||
PGW_LIMITER_BACKEND=memory # memory | redis(redis 需 REDIS_URL;多进程 worker 必须 redis)
|
PGW_LIMITER_BACKEND=memory # memory | redis(redis 需 REDIS_URL;多进程 worker 必须 redis)
|
||||||
|
|||||||
@@ -1,5 +1,43 @@
|
|||||||
# Changelog
|
# Changelog
|
||||||
|
|
||||||
|
## 1.3.6(2026-09-10)
|
||||||
|
|
||||||
|
给一次逻辑调用加了一条**可选**墙钟硬边界(issue #22),并修好取消路径的 TPM 结算与 `Retry-After` 非有限值防御。
|
||||||
|
|
||||||
|
### 关于调用期限,请先读这三句
|
||||||
|
|
||||||
|
| # | 承诺 | 展开 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 1 | **期限治理的是「等待」,不是「返回时刻」** | 到期后取消在飞的尝试,但清理(遥测写入、限流结算、缓存收尾)仍在 `finally` 里跑完,**允许超出期限**。实测构造(两次慢遥测写入)里返回时刻达期限的 **5–7 倍**;库只对「等待被切断」给承诺,对「多久返回」不给上界 |
|
||||||
|
| 2 | **到期 ≠ 未产出、≠ 未计费** | 上游可能已经算完并计费,只是结果在返回路上被丢弃(缓存写入慢于期限就是一例)。把 `CallDeadlineExceeded` 当成「这次没花钱」会低估成本 |
|
||||||
|
| 3 | **不配置就是 1.3.5 语义,逐字不变** | `call_deadline_s` 缺省 `None` 时根本不进 `asyncio.timeout` 上下文。故 1.3.5 的两条长等仍在: 纯 429 序列(429 不消耗重试预算)仍可能长时间等待;**有限大的 `Retry-After`(如 3600s)仍照睡**——库有意不用 `backoff_max_s` 去夹它,唯一制约手段就是本版这条期限 |
|
||||||
|
|
||||||
|
### 公共面新增(三项,全为纯新增)
|
||||||
|
|
||||||
|
| # | 位置 | 内容 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 1 | `polygateway.CallDeadlineExceeded` | 新异常;带 `scope` / `deadline_s`,**无 `retry_after_s`**(到期不含「何时可再试」,给 `0.0` 会指示下游立刻重打饱和渠道) |
|
||||||
|
| 2 | 配置键 `{SCOPE}__CALL_DEADLINE_S` | 缺省不设 = 不启用;非法值(0/负/`nan`/`inf`/非数)在**装配期**当场 `ValueError` |
|
||||||
|
| 3 | 三个 client 构造参数 + 四个公开方法的 keyword-only 参数 | `GatewayClient` / `EmbeddingClient` / `OcrClient` 的 `call_deadline_s`;`chat` / `embed` / `recognize_text` / `parse_layout` 可 per-call 覆盖(`None` = 继承装配值,**不提供「本次关闭」**)。一次 `embed` 的 N 个批次共享同一份期限,不随批数放大 |
|
||||||
|
|
||||||
|
> [!WARNING]
|
||||||
|
> **`except GatewayUnavailableError` 接不住 `CallDeadlineExceeded`。** 新异常直接继承 `PolyGatewayError`,既不属四分类,也不在 `GatewayUnavailableError` 族内——期限到期是**调用方自己设的边界**,不是网关不可用。只有显式配了期限的调用方才会遇到它,需要处理就单列一条 `except`。遥测侧无需改动: 三个边界既有的 `except PolyGatewayError` 会接住它并照常写一条 `terminal_failure` 行(`error_type='CallDeadlineExceeded'`),**零新增列**。
|
||||||
|
|
||||||
|
### 行为变更:取消路径的 TPM 结算口径
|
||||||
|
|
||||||
|
| 情形 | 1.3.5 | 本版 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 取消发生在**端口已开始、结算尚未确定**时 | `settle(0)`,入场预扣整笔退还 | 按 `est` **保留预扣**(方向是宁多扣不空退: 上游可能已计费) |
|
||||||
|
| 结算已确定(含真实 usage 恰为 0 的成功、已判 `SourceDead` 的 `0`) | 按已算出的值 | **一字不变**,取消不覆写 |
|
||||||
|
| 未被四分类接住的异常逃逸(`RuntimeError` 等) | `0` | **仍按 `0`**,本版不扩大语义(已登记为残留) |
|
||||||
|
|
||||||
|
启用期限后库自身会常规性触发取消路径,故这条记账修复与期限同版交付。OCR 的 `settle(0)` 不变——无 token 是事实而非「未知」。非取消路径的最终结算值与 1.3.5 逐字相同,只是算得更早(失败分支的结算决定前移到其第一个 `await` 之前)。
|
||||||
|
|
||||||
|
### 其他
|
||||||
|
|
||||||
|
- `Retry-After: inf` / `1e999`(`float()` 会把它舍成 `inf`)此前会原样进入退避并让该次尝试睡到天荒地老;现按「无提示」处理,退回纯指数退避,并发**一条**带源名与判据词的 warning(不回显原始头,429 风暴下会淹掉真信号)。`nan`、空串、负数、HTTP-date 的既有值语义一字未动。
|
||||||
|
- 限流 Lua、`Permit` 端口签名、遥测 schema、缓存 key 公式、重试预算与退避算法、熔断语义**均未改动**。
|
||||||
|
|
||||||
## 1.3.5(2026-09-09)
|
## 1.3.5(2026-09-09)
|
||||||
|
|
||||||
把治理单位从「一次尝试」补齐到「一次逻辑调用」(issue #19、#23)。此前重试、换源、结构化重问、embedding 分批都各自独立可见,而「这一次调用总共打了几次、总共花了多久、最后为什么失败」在库外拼不出来;结构化耗尽、embedding/OCR 的无源与准入拒绝更是**一条遥测行都没有**。
|
把治理单位从「一次尝试」补齐到「一次逻辑调用」(issue #19、#23)。此前重试、换源、结构化重问、embedding 分批都各自独立可见,而「这一次调用总共打了几次、总共花了多久、最后为什么失败」在库外拼不出来;结构化耗尽、embedding/OCR 的无源与准入拒绝更是**一条遥测行都没有**。
|
||||||
|
|||||||
@@ -91,7 +91,7 @@ make ci # 只读验证(check + test)
|
|||||||
| 1 | **更新 README** | 打包会把当时的 README 固化进 sdist,**发布后再改就来不及了**(包里那份永远是旧的)。逐项核对: 安装命令的版本约束(`==1.1.*` 这类**极易漏改**,漏了下游就被锁在旧版)、能力表是否覆盖新行为、数字型断言是否仍成立(如遥测字段数,须用 `inspect.signature` 实测而非凭记忆) |
|
| 1 | **更新 README** | 打包会把当时的 README 固化进 sdist,**发布后再改就来不及了**(包里那份永远是旧的)。逐项核对: 安装命令的版本约束(`==1.1.*` 这类**极易漏改**,漏了下游就被锁在旧版)、能力表是否覆盖新行为、数字型断言是否仍成立(如遥测字段数,须用 `inspect.signature` 实测而非凭记忆) |
|
||||||
| 2 | CHANGELOG 定版 | "未发布" → `## X.Y.Z(日期)` |
|
| 2 | CHANGELOG 定版 | "未发布" → `## X.Y.Z(日期)` |
|
||||||
| 3 | 版本号 | `pyproject.toml` + `src/polygateway/__init__.py` 两处必须一致 |
|
| 3 | 版本号 | `pyproject.toml` + `src/polygateway/__init__.py` 两处必须一致 |
|
||||||
| 4 | 合并 main + push | `--no-ff`;合并后在 main 上重跑 `make lint` 与全套件,**外加 `pytest -m slow`** ——真实网关 e2e 与 Redis 时间语义变体被 `addopts = "-m 'not slow'"` 默认排除,**不显式跑就等于没跑**(约 20-40 分钟,取决于网关快慢)。它们不进日常提交是有意的: pre-commit 关卡跑全套件,网关一抖就挡住与之无关的提交,久了会把"测试红了先怀疑网关"变成惯性,真 bug 也会被当成抖动重试掉;代价是这道门必须由本清单兜住 |
|
| 4 | 合并 main + push | `--no-ff`;合并后在 main 上重跑 `make lint` 与全套件,**外加与本次 diff 有交集的 slow 子集**(2026-09-10 起): `pytest -m slow` 只跑被本版改动触及的模块——Redis 时间语义变体在动限流/退避/取消时跑,真实网关冒烟在动公开入口时跑;**`test_thinking_live.py` 全模型能力矩阵不再每次发布都跑**(付费且额度耗尽渠道会走完整超时链,一晚数小时无信号),仅在动 `thinking.py`/能力注册表/相关 e2e 设施或人类明确要求时跑,否则复用最近一次有效矩阵证据并如实登记。slow 不进日常提交是有意的: pre-commit 关卡跑全套件,网关一抖就挡住与之无关的提交,久了会把"测试红了先怀疑网关"变成惯性,真 bug 也会被当成抖动重试掉;代价是这道门必须由本清单兜住 |
|
||||||
| 5 | **打 tag 并 push** | `git tag -a vX.Y.Z -m "..."` + `git push origin vX.Y.Z`。历史上多个版本漏打 |
|
| 5 | **打 tag 并 push** | `git tag -a vX.Y.Z -m "..."` + `git push origin vX.Y.Z`。历史上多个版本漏打 |
|
||||||
| 6 | 构建 | `rm -rf dist && python -m build && python -m twine check dist/*` |
|
| 6 | 构建 | `rm -rf dist && python -m build && python -m twine check dist/*` |
|
||||||
| 7 | **上传 registry** | 凭据在 `~/.config/tea/config.yml`(tea CLI 的 Gitea token,**不在** `~/.pypirc`);token 走 `TWINE_PASSWORD` 环境变量,不进命令行<br>`TWINE_USERNAME=iomgaa TWINE_PASSWORD=$TOKEN python -m twine upload --repository-url https://gitea.iomgaa.online/api/packages/iomgaa/pypi dist/*` |
|
| 7 | **上传 registry** | 凭据在 `~/.config/tea/config.yml`(tea CLI 的 Gitea token,**不在** `~/.pypirc`);token 走 `TWINE_PASSWORD` 环境变量,不进命令行<br>`TWINE_USERNAME=iomgaa TWINE_PASSWORD=$TOKEN python -m twine upload --repository-url https://gitea.iomgaa.online/api/packages/iomgaa/pypi dist/*` |
|
||||||
@@ -115,7 +115,7 @@ Gitea 包 registry 是 **owner 级**(`/iomgaa/-/packages/`)不是仓库级;PyPI
|
|||||||
- 覆盖率目标 80%;并发/韧性行为是一等测试对象: 重试穿透取消、熔断开路半开、限流结算退款、Redis 掉线降级方向、缓存 key 隔离。
|
- 覆盖率目标 80%;并发/韧性行为是一等测试对象: 重试穿透取消、熔断开路半开、限流结算退款、Redis 掉线降级方向、缓存 key 隔离。
|
||||||
- Redis 相关测试用真实 Redis(integration),不 mock Lua 行为;限流契约测试随实现一起交付(参考 CHSAnalyzer `tests/contracts_limiter.py`)。
|
- Redis 相关测试用真实 Redis(integration),不 mock Lua 行为;限流契约测试随实现一起交付(参考 CHSAnalyzer `tests/contracts_limiter.py`)。
|
||||||
- 涉及真实 LLM 的测试输出结构化 Markdown 至 `tests/outputs/<module>/<test>_<ts>.md`。
|
- 涉及真实 LLM 的测试输出结构化 Markdown 至 `tests/outputs/<module>/<test>_<ts>.md`。
|
||||||
- **成败取决于外部服务当下状态的测试一律标 `slow`**(`tests/e2e/` 四个文件与 Redis 时间语义变体):它们默认不进日常套件,由发布清单第 4 步统一跑。判据是"重跑一次可能就绿了"——这种测试留在提交关卡里会污染信号。同理,给它们的超时不得紧于 `.env` 的生产配置,否则是设计上就会间歇红。
|
- **成败取决于外部服务当下状态的测试一律标 `slow`**(`tests/e2e/` 四个文件与 Redis 时间语义变体):它们默认不进日常套件,由发布清单第 4 步**按 diff 交集选子集**跑(全量 `pytest -m slow` 仅在交集不清或人类要求时用)。判据是"重跑一次可能就绿了"——这种测试留在提交关卡里会污染信号。同理,给它们的超时不得紧于 `.env` 的生产配置,否则是设计上就会间歇红。
|
||||||
|
|
||||||
## 5. 项目结构
|
## 5. 项目结构
|
||||||
|
|
||||||
|
|||||||
@@ -16,6 +16,7 @@
|
|||||||
| 熔断 | 双通道(连续失败 + 失败率窗口,健康证据抑制误熔);半开单探针带租约(持有者死亡自动回收);epoch fencing 拒绝迟到写回;开路时长指数递增;**开路时当场失败还是等冷却可配**(`CIRCUIT_OPEN`,单源 scope 应配 `wait`) |
|
| 熔断 | 双通道(连续失败 + 失败率窗口,健康证据抑制误熔);半开单探针带租约(持有者死亡自动回收);epoch fencing 拒绝迟到写回;开路时长指数递增;**开路时当场失败还是等冷却可配**(`CIRCUIT_OPEN`,单源 scope 应配 `wait`) |
|
||||||
| 自适应并发 | AIMD:429 削减、成功缓升,防止打爆上游 |
|
| 自适应并发 | AIMD:429 削减、成功缓升,防止打爆上游 |
|
||||||
| 背压与判死 | 配额满与熔断开路**各自**可选等待或快速失败(`QUOTA_FULL` / `CIRCUIT_OPEN`,两键不可互相替代);等待期按双条件判死(本地非生产性等待与全局无进展**同时**超窗)。stall 窗口只计**非生产性**等待(429 退避/配额轮询/熔断冷却),与 `TIMEOUT_S` 无耦合 |
|
| 背压与判死 | 配额满与熔断开路**各自**可选等待或快速失败(`QUOTA_FULL` / `CIRCUIT_OPEN`,两键不可互相替代);等待期按双条件判死(本地非生产性等待与全局无进展**同时**超窗)。stall 窗口只计**非生产性**等待(429 退避/配额轮询/熔断冷却),与 `TIMEOUT_S` 无耦合 |
|
||||||
|
| 调用期限 | 一次逻辑调用可选一条**墙钟硬边界**(`{SCOPE}__CALL_DEADLINE_S` 或 `chat(call_deadline_s=...)`,缺省不启用):治理的是**等待**——重试退避、配额轮询、熔断冷却、结构化重问与 embedding 分批共享同一份期限。三条须知:①**返回时刻 = 期限 + 清理耗时**(遥测/结算/缓存收尾在 `finally` 里跑完,允许超期;实测构造达期限的 5–7 倍),库只承诺切断等待、不给返回上界;②**到期 ≠ 未产出、≠ 未计费**,上游可能已算完并计费;③**不配置即逐字保持 1.3.5 语义**(纯 429 序列仍可能长等、有限大 `Retry-After` 仍照睡)。到期抛 `CallDeadlineExceeded`,**不属四分类、不属 `GatewayUnavailableError` 族** |
|
||||||
| 响应缓存 | Redis/内存;key 含 model + messages 摘要 + namespace(缓存隔离单位)+ salt + 采样参数 + 请求级推理档位(同 messages 跑 low 与 max 不互相命中),多模态 content 先摘要再 hash(防毒化);可 per-call 绕过(科研重采样) |
|
| 响应缓存 | Redis/内存;key 含 model + messages 摘要 + namespace(缓存隔离单位)+ salt + 采样参数 + 请求级推理档位(同 messages 跑 low 与 max 不互相命中),多模态 content 先摘要再 hash(防毒化);可 per-call 绕过(科研重采样) |
|
||||||
| 流式看门狗 | TTFT / inter-token / 总超时三层活性;thinking token 刷活性不计结果;截断流(缺 `[DONE]`)判瞬时不入缓存 |
|
| 流式看门狗 | TTFT / inter-token / 总超时三层活性;thinking token 刷活性不计结果;截断流(缺 `[DONE]`)判瞬时不入缓存 |
|
||||||
| 推理可观测性 | "这次到底推理没推理"由多信号裁定(推理正文压倒 usage 明细),三态落在 `LLMResponse.thinking_observation`:`observed` / `absent` / `unknown`——**`unknown` 是"本次判不出",不是"没推理"**;本次实发档位与实测观测矛盾时按 `(源, 模型, 生效档位)` 各告警一次(能力表过期、开启未生效、注入了却观测不到;同一模型的 low 与 max 是两个独立的矛盾,不共用节流键);裁定结果随遥测落库 |
|
| 推理可观测性 | "这次到底推理没推理"由多信号裁定(推理正文压倒 usage 明细),三态落在 `LLMResponse.thinking_observation`:`observed` / `absent` / `unknown`——**`unknown` 是"本次判不出",不是"没推理"**;本次实发档位与实测观测矛盾时按 `(源, 模型, 生效档位)` 各告警一次(能力表过期、开启未生效、注入了却观测不到;同一模型的 low 与 max 是两个独立的矛盾,不共用节流键);裁定结果随遥测落库 |
|
||||||
@@ -115,7 +116,7 @@ stats.total_latency_ms # 含缓存 IO、退避、准入等待、重
|
|||||||
|
|
||||||
```bash
|
```bash
|
||||||
pip install --extra-index-url https://gitea.iomgaa.online/api/packages/iomgaa/pypi/simple/ \
|
pip install --extra-index-url https://gitea.iomgaa.online/api/packages/iomgaa/pypi/simple/ \
|
||||||
"polygateway[redis,postgres,structured]>=1.3.5,<2"
|
"polygateway[redis,postgres,structured]>=1.3.6,<2"
|
||||||
```
|
```
|
||||||
|
|
||||||
核心仅依赖 `httpx` + `pydantic`;按需选 extras:
|
核心仅依赖 `httpx` + `pydantic`;按需选 extras:
|
||||||
@@ -185,13 +186,17 @@ vectors = (await embed.embed(["文本 a", "文本 b"])).vectors
|
|||||||
### 4. 业务侧异常处理
|
### 4. 业务侧异常处理
|
||||||
|
|
||||||
```python
|
```python
|
||||||
from polygateway import GatewayUnavailableError, RequestRejectedError
|
from polygateway import CallDeadlineExceeded, GatewayUnavailableError, RequestRejectedError
|
||||||
|
|
||||||
try:
|
try:
|
||||||
resp = await client.chat(messages)
|
resp = await client.chat(messages)
|
||||||
except GatewayUnavailableError as exc:
|
except GatewayUnavailableError as exc:
|
||||||
# 整个 scope 暂时无源可用: 延期重投,不消耗业务失败预算
|
# 整个 scope 暂时无源可用: 延期重投,不消耗业务失败预算
|
||||||
schedule_retry(after_s=exc.retry_after_s) # exc.reason / exc.per_source_reasons 供诊断
|
schedule_retry(after_s=exc.retry_after_s) # exc.reason / exc.per_source_reasons 供诊断
|
||||||
|
except CallDeadlineExceeded as exc:
|
||||||
|
# 只在自己配了调用期限时出现: 不在 GatewayUnavailableError 族内,上一条接不住;
|
||||||
|
# 且无 retry_after_s(到期不含"何时可再试"),重投时机由业务侧定
|
||||||
|
schedule_retry(after_s=None) # exc.scope / exc.deadline_s 供诊断
|
||||||
except RequestRejectedError:
|
except RequestRejectedError:
|
||||||
... # 请求本身有问题(400/格式拒绝): 不重试,直接失败
|
... # 请求本身有问题(400/格式拒绝): 不重试,直接失败
|
||||||
```
|
```
|
||||||
@@ -463,6 +468,8 @@ SQLite 侧**不建议**对着一个大库文件跑 `DELETE` + `VACUUM`,而应**
|
|||||||
|
|
||||||
**网关拒绝的理由不会丢失**(1.2.0 起):非 2xx 的响应体经折叠与截断后同时进入异常 message 与 `exc.body_text`,故遥测表的 `error` 列里就能看到网关的原话——不必再为查一次 400 单独埋点。截断保头保尾(总长 2048 字符),JSON 错误体尾部的 `code` / `request_id` 不会被切掉。**经中转部署时请注意**:第三方中转服务自身抖动也会回 400,从状态码上与"你的输入有问题"无法区分;库仍按确定性失败处理(直连供应商时重试只会白烧配额),批处理下游宜据 `body_text` 自备兜底分类。
|
**网关拒绝的理由不会丢失**(1.2.0 起):非 2xx 的响应体经折叠与截断后同时进入异常 message 与 `exc.body_text`,故遥测表的 `error` 列里就能看到网关的原话——不必再为查一次 400 单独埋点。截断保头保尾(总长 2048 字符),JSON 错误体尾部的 `code` / `request_id` 不会被切掉。**经中转部署时请注意**:第三方中转服务自身抖动也会回 400,从状态码上与"你的输入有问题"无法区分;库仍按确定性失败处理(直连供应商时重试只会白烧配额),批处理下游宜据 `body_text` 自备兜底分类。
|
||||||
|
|
||||||
|
**有限大的 `Retry-After` 仍照睡**: 能力表那条「尊重 `Retry-After`」是字面意思——库**有意不用 `backoff_max_s` 去夹服务端给的提示**(夹住就是提前重打已明确说「还没好」的网关),服务端给 3600s 就真睡 3600s;**唯一的制约手段是 1.3.6 的调用期限**(上表「调用期限」行)。自 1.3.6 起,`inf` / `-inf` / `1e999` 这类**非有限**取值按「无提示」处理(退回纯指数退避 + 一条带源名的 warning),不再造成无限等待;`nan`、空串、负数、HTTP-date 的既有语义不变。
|
||||||
|
|
||||||
### 哪些异常会到达调用方
|
### 哪些异常会到达调用方
|
||||||
|
|
||||||
上表的"库内行为"一列描述的是**治理动作**,不是调用方要处理的东西。四类里有两类**根本到不了调用方**——它们被重试循环接住,预算耗尽时统一包成 `AllSourcesExhausted`。这个区分只看类型树和 docstring 是读不出来的,曾让下游据此写错整段设计文档,故在此列明:
|
上表的"库内行为"一列描述的是**治理动作**,不是调用方要处理的东西。四类里有两类**根本到不了调用方**——它们被重试循环接住,预算耗尽时统一包成 `AllSourcesExhausted`。这个区分只看类型树和 docstring 是读不出来的,曾让下游据此写错整段设计文档,故在此列明:
|
||||||
@@ -473,9 +480,12 @@ SQLite 侧**不建议**对着一个大库文件跑 `DELETE` + `VACUUM`,而应**
|
|||||||
| `RequestRejectedError` | `SourceDeadError`(立即熔断该源并换源,同上) |
|
| `RequestRejectedError` | `SourceDeadError`(立即熔断该源并换源,同上) |
|
||||||
| `ResultInvalidError` | |
|
| `ResultInvalidError` | |
|
||||||
| `SourceNotConfiguredError` | |
|
| `SourceNotConfiguredError` | |
|
||||||
|
| `CallDeadlineExceeded`(1.3.6 起) | |
|
||||||
|
|
||||||
**`GovernanceBackendError` 属于第一列**: 限流/熔断的状态后端(如 Redis)自身故障时库 fail-closed——一个请求都发不出去,这就是"整个 scope 暂时不可用"。它继承 `GatewayUnavailableError`,所以 §4 那段 `except GatewayUnavailableError` 一条即覆盖完整,无需为它单列分支。`retry_after_s` 默认 5 秒(后端恢复时间不可知,取 0 会让积压任务零延迟冲击已挂掉的后端)。
|
**`GovernanceBackendError` 属于第一列**: 限流/熔断的状态后端(如 Redis)自身故障时库 fail-closed——一个请求都发不出去,这就是"整个 scope 暂时不可用"。它继承 `GatewayUnavailableError`,所以 §4 那段 `except GatewayUnavailableError` 一条即覆盖完整,无需为它单列分支。`retry_after_s` 默认 5 秒(后端恢复时间不可知,取 0 会让积压任务零延迟冲击已挂掉的后端)。
|
||||||
|
|
||||||
|
**`CallDeadlineExceeded` 既不属四分类、也不属 `GatewayUnavailableError` 族**: 它直接继承 `PolyGatewayError`,表达的是「调用方自己设的墙钟边界到了」而非网关不可用,故**只有显式配了 `call_deadline_s` / `{SCOPE}__CALL_DEADLINE_S` 的调用方才可能遇到它**,不配就永远不会出现。它**无 `retry_after_s`**,且 §4 那段 `except GatewayUnavailableError` **接不住**它——要处理就得单列一条分支。
|
||||||
|
|
||||||
**`SourceNotConfiguredError` 有意不在第一列的族内**: 源名不在限流后端的配置字典中是**装配缺陷**而非暂时故障,它应当消耗失败预算、进死信、让人看见——归入可重投家族只会让配置写错的任务永远重投且无人告警。
|
**`SourceNotConfiguredError` 有意不在第一列的族内**: 源名不在限流后端的配置字典中是**装配缺陷**而非暂时故障,它应当消耗失败预算、进死信、让人看见——归入可重投家族只会让配置写错的任务永远重投且无人告警。
|
||||||
|
|
||||||
## 配置参考
|
## 配置参考
|
||||||
|
|||||||
+2
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "polygateway"
|
name = "polygateway"
|
||||||
version = "1.3.5"
|
version = "1.3.6"
|
||||||
description = "PolyGateway:实验室统一的大语言模型(LLM/VLM/OCR)调度与中转库——多源、限流、重试、熔断、缓存、遥测"
|
description = "PolyGateway:实验室统一的大语言模型(LLM/VLM/OCR)调度与中转库——多源、限流、重试、熔断、缓存、遥测"
|
||||||
# registry 包页面的正文只认这一项:缺了页面就是一片空白(1.1.2 的教训,twine 会警告
|
# registry 包页面的正文只认这一项:缺了页面就是一片空白(1.1.2 的教训,twine 会警告
|
||||||
# long_description missing 但不阻塞上传)。README 在打包时被固化进产物,发布后再改无效。
|
# long_description missing 但不阻塞上传)。README 在打包时被固化进产物,发布后再改无效。
|
||||||
@@ -82,6 +82,7 @@ layers = [
|
|||||||
"polygateway.middleware",
|
"polygateway.middleware",
|
||||||
"polygateway.transports | polygateway.backends | polygateway.telemetry | polygateway.structured",
|
"polygateway.transports | polygateway.backends | polygateway.telemetry | polygateway.structured",
|
||||||
"polygateway.thinking",
|
"polygateway.thinking",
|
||||||
|
"polygateway.deadline",
|
||||||
"polygateway.providers : polygateway.sources",
|
"polygateway.providers : polygateway.sources",
|
||||||
"polygateway.ports : polygateway.types : polygateway.errors : polygateway.streaming",
|
"polygateway.ports : polygateway.types : polygateway.errors : polygateway.streaming",
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -0,0 +1,313 @@
|
|||||||
|
# 1.3.6:可选调用期限(issue #22),兼 `Retry-After` 有限值防御
|
||||||
|
|
||||||
|
- 状态: **人类已于 2026-09-10 批准**(§9 七项批准项全数获批,H3 取 (a′) 档);实施计划见 `research-wiki/plans/2026-09-10-136-call-deadline.md`
|
||||||
|
- 基线: main `ab00aa4` / 1.3.5,工作区 HEAD `d2455e8`;分支 `feature/1.3.6-call-budgets`
|
||||||
|
- 范围收敛(2026-09-09 人类定夺,2026-09-10 追认): **只做方案 B 的后半**——保留 chat 的 429 免次数预算,新增**可选**整体调用期限,缺省不启用;**不做** 429 次数预算(原 D2/D3)、**不做** issue #24(未批,不得捆绑)
|
||||||
|
- 已批准契约(2026-09-10 逐条定案): `call_deadline_s` 缺省 `None`;单调用参数 `None` = 继承装配值;`CallDeadlineExceeded` 为 `PolyGatewayError` **直接**子类;合作式清理会超出期限、**到期可能丢弃已计费的成功**;chat 429 免次数预算保留原样
|
||||||
|
- 取消结算(TPM)修复**并入本版**: 不再作为「先立独立 issue、修好再启用期限」的前置阻塞,改为本版第一个原子提交(§6.3 矩阵即其精确契约)
|
||||||
|
- 输入: issue #22 原文;设计审查 `a544b789/design136/review.md`(BL1-BL5)与独立复审 B(B1-B5,含离线 asyncio 探针实测);本轮独立探针 `/tmp/pgw_deadline_probe.py`、`/tmp/pgw_deadline_probe2.py`(3.12.13 实测,§5.2);1.3.5 源码逐行现读
|
||||||
|
- 关联: ARCHITECTURE §7.2/§7.3、`designs/2026-08-06-issue8-stall-budget-design.md`、`designs/2026-09-09-135-call-observability-design.md`
|
||||||
|
|
||||||
|
## 1. 目标与非目标
|
||||||
|
|
||||||
|
| 项 | 内容 |
|
||||||
|
| --- | --- |
|
||||||
|
| 目标 1 | 让调用方能对**一次逻辑调用**设墙钟上限;不配置时,库行为逐字保持 1.3.5 |
|
||||||
|
| 目标 2 | 修 `Retry-After` **非有限值**防御缺口(`inf`/`1e999` → `sleep(inf)` 永久挂起) |
|
||||||
|
| 目标 3 | 期限是**单一硬边界**,覆盖缓存 IO/准入排队/退避 sleep/transport/结构化重问/embedding 分批,不是轮首软检查 |
|
||||||
|
| 非目标 A | 不新增 429 次数预算、不改 429 退避指数分账、不收紧任何缺省(原 §4.1/§4.2 已删除) |
|
||||||
|
| 非目标 B | 不改三条循环各自的既有差异(chat 429 免预算 + stall 退还;embed/OCR 无条件计数) |
|
||||||
|
| 非目标 C | ~~本设计不含取消路径 `settle(0)` 结算修复~~ → **2026-09-10 改判**: 该修复**已获批并并入本版**,是期限落地前的第一个原子提交(§6.3 给出精确结算矩阵)。它只改**取消路径**的结算取值,不动成功/拒绝/失败三条既有分支的口径 |
|
||||||
|
| 非目标 D | 不做 issue #24(长尾对冲);不新增遥测列;不新增后台任务/`shield` |
|
||||||
|
|
||||||
|
## 2. 1.3.5 现状核实(现读源码,不引用旧报告)
|
||||||
|
|
||||||
|
| 事实 | 证据 | 对本设计的意义 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| chat 429 免重试预算并退还 stall 账 | `middleware/retry.py:158-164`、`:250-257` | **保留不动**;期限是正交的第二层 |
|
||||||
|
| embed/OCR 无条件 `fails += 1` | `embedding.py:319`、`ocr.py:339` | **保留不动**;两条循环次数有界、时长无界,由期限兜 |
|
||||||
|
| 退避与 `Retry-After` 取大且不夹上限 | `retry.py:64-77` `max(delay, retry_after)` | 有限大值**照睡**(§6.2 说明为何不夹) |
|
||||||
|
| `_parse_retry_after` 放行 `inf`,但 `nan` **已被忽略** | `transports/openai_compat.py:111-120`:`float()` 成功后 `seconds > 0`;`nan > 0` 恒假 → `nan` 现已返回 `None` | 真实缺口只有 `inf`/`1e999` 一族(唯一一处;`monkey_ocr.py:58` 不解析该头) |
|
||||||
|
| 三个公开边界均在**输入校验之后**建 `_CallContext` | `client.py:381`、`embedding.py:192`、`ocr.py:272` | 期限起点与 `total_latency_ms` 同口径,校验时间不计入 |
|
||||||
|
| 洋葱与两条循环全在同一 `await` 树下 | `client.py:398`、`embedding.py:209`、`ocr.py:274` | 一个 `asyncio.timeout` 即可覆盖全部等待,**无须**改 `backoff_delay`/`SourceAdmission._nap` 签名(解 BL3) |
|
||||||
|
| 终态行唯一出口 + 去重 | `telemetry.py:670-716` `emit_terminal_once` / `types.py:355` `claim_terminal` | 到期终态复用该出口,`error_type` 列自动落新类名,**无新增列** |
|
||||||
|
| 取消路径 attempt 行写 `error="cancelled"` 后穿透 | `retry.py:320-323` | 到期时 attempt 行仍记 cancelled(保留),**终态行**必须记 deadline |
|
||||||
|
| 库内已有 `asyncio.timeout` + `cm.expired()` 范式 | `streaming.py:43-58`、`telemetry/postgres.py:419-423` | 本设计沿用同一范式,不发明新写法 |
|
||||||
|
| `asyncio.timeout` 的**公共契约**: 到期投递取消并在退出时转 `TimeoutError`,外部取消原样上抛,擦边成功不遗留游离取消 | 标准库文档 + 受支持 Python 矩阵上的行为测试(§10"形态区分"批次;复审 B 在 3.12.13 上以离线探针复核) | 外部取消优先与"擦边成功"竞态由运行时判定,库不自判(§5.2);**不依据 CPython 私有实现细节立论**,跨版本保证由测试矩阵给 |
|
||||||
|
|
||||||
|
## 3. 实现范式的三个备选
|
||||||
|
|
||||||
|
| 方案 | 做法 | 判定 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| **A 单一硬边界(推荐)** | 在三个公开边界各用一次 `asyncio.timeout(d)` 包住整条 `await` 树;到期由运行时取消在途 await,边界处转成 `CallDeadlineExceeded` | 覆盖面天然完整(缓存/准入/sleep/transport/重问/分批);零签名扩散;代价是**到期即取消在途尝试**(§6.3) |
|
||||||
|
| B 轮首软检查 + sleep 夹紧 | 三条循环轮首判剩余预算,退避 `min(delay, 剩余)` | **否决**: 名为期限实为"轮次粒度上限"——一次 300s 的慢尝试或一次准入排队即可整体越限,且要改 `backoff_delay`、`_nap`、三处轮首共 5 个点(BL3 原形态)。用软检查冒充硬期限正是 issue #8 "不要用一个预算冒充另一个预算"的同形错误 |
|
||||||
|
| C 每层各自超时 | 缓存、准入、transport 各配一份超时 | **否决**: N 个键相加才是总时长,调用方仍拿不到"最多等 N 秒"的承诺;配置面爆炸 |
|
||||||
|
|
||||||
|
**推荐 A**,且缺省 `None` = 不启用:opt-in 才能保证存量下游行为逐字不变(1.3.x 内不做默认行为变更)。
|
||||||
|
|
||||||
|
## 4. 方案 A 的具体形态
|
||||||
|
|
||||||
|
### 4.1 新模块 `src/polygateway/deadline.py`(约 50 行,只依赖 stdlib + `errors.py`)
|
||||||
|
|
||||||
|
```python
|
||||||
|
def ensure_call_deadline(value: object, origin: str) -> float | None:
|
||||||
|
"""全装配路径共用的值域校验: None 或有限正数,否则 ValueError(消息含 origin)。"""
|
||||||
|
if value is None:
|
||||||
|
return None
|
||||||
|
if isinstance(value, bool) or not isinstance(value, (int, float)):
|
||||||
|
raise ValueError(...) # bool 先判:不得把 True 当成 1 秒
|
||||||
|
v = float(value)
|
||||||
|
if not math.isfinite(v) or v <= 0:
|
||||||
|
raise ValueError(...) # NaN / inf / 0 / 负
|
||||||
|
return v
|
||||||
|
|
||||||
|
|
||||||
|
async def with_call_deadline[T](aw: Awaitable[T], *, deadline_s: float | None, scope: str) -> T:
|
||||||
|
"""给一次逻辑调用施加单一硬边界;None = 不进任何上下文,逐字走旧路径。
|
||||||
|
|
||||||
|
入参已由调用方在**构造 `aw` 之前**校验(见 §4.3),本函数不再校验。
|
||||||
|
"""
|
||||||
|
if deadline_s is None:
|
||||||
|
return await aw
|
||||||
|
inner_timeout: BaseException | None = None
|
||||||
|
try:
|
||||||
|
async with asyncio.timeout(deadline_s) as cm:
|
||||||
|
try:
|
||||||
|
return await aw
|
||||||
|
except TimeoutError as exc:
|
||||||
|
inner_timeout = exc # 体内(含清理路径)自抛,不是本层期限
|
||||||
|
raise
|
||||||
|
except TimeoutError as exc:
|
||||||
|
if cm.expired() and exc is not inner_timeout:
|
||||||
|
raise CallDeadlineExceeded(scope=scope, deadline_s=deadline_s) from None
|
||||||
|
raise # 内层失败原样上抛,绝不贴 deadline 标签
|
||||||
|
```
|
||||||
|
|
||||||
|
**为何不能只用 `cm.expired()`(本轮独立探针实测,3.12.13)**: `Timeout.expired()` 在 `EXPIRING`/`EXPIRED` 两态都返 True(`inspect.getsource` 现读),而计时器触发后的**清理路径**若自抛 `TimeoutError`(第三方端口实现自抛,或其内部另一个 `asyncio.timeout` 到期),该异常会被只看 `expired()` 的写法**改标成 `CallDeadlineExceeded`**(探针 E1/E2 实测均为误标)。`__cause__` 启发式也不够: 内层 `asyncio.timeout` 抛的 `TimeoutError` 其 `__cause__` 同样是 `CancelledError`(E1 实测仍误标)。故采用**局部变量身份比较**这一最小辅助机制: 本层 `asyncio.timeout` 转出的是 `raise TimeoutError from exc_val` 新建对象,与体内那个实例必不同一,判据确定、不依赖任何 CPython 私有实现。它**不新增公共配置、不开后台任务、不改异常对象**;`cm.expired()` 作为第二道守卫保留。
|
||||||
|
|
||||||
|
```text
|
||||||
|
探针证据(/tmp/pgw_deadline_probe2.py, python 3.12.13):
|
||||||
|
D1 到期命中 → CallDeadlineExceeded 耗时 0.050s cancelling=0
|
||||||
|
D2 清理内层 timeout → TimeoutError(原样) 耗时 0.060s ← 只看 expired() 会误标
|
||||||
|
D3 清理裸 TimeoutError→ TimeoutError(原样) 耗时 0.050s ← 只看 expired() 会误标
|
||||||
|
D4 未到期内层自抛 → TimeoutError(原样) 耗时 0.010s
|
||||||
|
D5 到期窗口内领域异常 → 领域异常(期限让位) 耗时 0.200s
|
||||||
|
D6 慢清理 → CallDeadlineExceeded 耗时 0.251s(期限 5×)
|
||||||
|
D7 清理期外部取消 → CancelledError cancelling=1(外部取消优先)
|
||||||
|
D8 外部取消先到 → CancelledError
|
||||||
|
D9 未启用(None) → 逐字旧路径
|
||||||
|
```
|
||||||
|
|
||||||
|
写成**接收 awaitable 的函数**而非 `@asynccontextmanager`: 后者要把 `yield` 包进 timeout,取消经 `athrow` 回注生成器,语义正确但绕(`streaming.py:60-70` 的 docstring 已记录同类陷阱);函数形态只有一条直路。`None` 分支**不进** `asyncio.timeout`,故未启用时连"必须在 Task 内运行"这一新约束都不引入。
|
||||||
|
|
||||||
|
**两个实现约束(实施时不得变形)**:
|
||||||
|
|
||||||
|
1. **校验先于构造 awaitable**: 公开方法入口先跑 `ensure_call_deadline`,通过后才构造 `self._handler(request)` 等协程。若把校验放进 `with_call_deadline`,非法参数抛错时会遗留**未 await 的协程**(RuntimeWarning + 未释放资源)。
|
||||||
|
2. **分层合法**: `deadline.py` 只 import stdlib 与 `errors.py`,并作为**新一层**进 import-linter 契约(`pyproject.toml` layers 中置于 `polygateway.thinking` 之下、内核行之上)。`config.py` 位在更上层,import 它合法且**不会依赖任何具体实现**(transports/backends/telemetry 一律不引入)。
|
||||||
|
|
||||||
|
### 4.2 三处接入点(唯一三处;严禁在循环内层再建 scope)
|
||||||
|
|
||||||
|
| 文件:行 | 现状 | 改后 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| `client.py:398` | `response = await self._handler(request)` | `await with_call_deadline(self._handler(request), deadline_s=d, scope=self._scope)` |
|
||||||
|
| `embedding.py:209` | `return await self._embed_all(...)` | 同款包住 `_embed_all(...)`(**整次调用一份**,分批共享) |
|
||||||
|
| `ocr.py:274` | `return await self._run(...)` | 同款包住 `_run(...)` |
|
||||||
|
|
||||||
|
三处均在既有 `try` 之内、`_CallContext` 之后,故到期路径照走 `except PolyGatewayError → emit_terminal_once`(§7)。`StructuredMW._run_ladder`(`structured.py:67-98`)与 `_embed_batch`(`embedding.py:295`)**不得**新建 scope:同级重试/重问/分批共享同一期限,否则期限被轮数放大 N 倍即等于没有。
|
||||||
|
|
||||||
|
### 4.3 装配路径与 per-call 覆盖(实际签名)
|
||||||
|
|
||||||
|
| 层 | 签名变化 | 语义 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 配置键 | `{SCOPE}__CALL_DEADLINE_S`,经 `_first` 读(**不用** `_require`,否则是破坏性配置变更) | 未设 = `None` = 不启用 |
|
||||||
|
| `GatewaySettings` | 末尾追加 `call_deadline_s: float \| None = None` + `_validate_call_deadline()` 进 `__post_init__`(见 `config.py:186-193`) | 有默认值,不扰动既有位置构造;校验覆盖 env/直接构造/`dataclasses.replace` 三条路 |
|
||||||
|
| 值域校验 | **全装配路径共用** `deadline.ensure_call_deadline`(§4.1):`GatewaySettings.__post_init__`、三个 client `__init__`、四个公开方法各调一次,实现只一份 | 拒收: `bool`(True 不得当 1 秒)、非数值类型、`NaN`、`inf`、`0`、负数 → `ValueError`(消息含来源)。**不与 `timeout_s` 耦合**:期限短于单次超时是合法选择 |
|
||||||
|
| 三个 client `__init__` | 追加 keyword-only `call_deadline_s: float \| None = None`,**入口即校** | 与 `now`/`sleep`/`rng` 同款注入位;全量注入是正式装配路,不得只靠 `GatewaySettings` 守门(否则 `inf` 静默失效、`NaN` 每次调用当场失败) |
|
||||||
|
| 三条 `from_settings` | 传 `settings.call_deadline_s`(embed/OCR 取 `settings.gateway.call_deadline_s`) | `from_env` 无签名变化(经 settings 透传) |
|
||||||
|
| 四个公开方法 | `chat`/`embed`/`recognize_text`/`parse_layout` 追加 keyword-only `call_deadline_s: float \| None = None` | `None` = **继承装配值**;正数 = 本次覆盖;**不提供"本次关闭"**(需要不同期限就装配两个 client;三态哨兵不值这个公共面复杂度) |
|
||||||
|
| 校验时点 | per-call 值在 `_CallContext` 创建**之前**、也在构造被包裹协程之前校验(与既有三项校验同列;OCR 两个入口经 `_call` 两级透传,与 `image` 校验同列) | 输入校验边界保持:非法期限抛裸 `ValueError`,不进统计边界、不写终态行、不遗留未 await 协程 |
|
||||||
|
|
||||||
|
### 4.4 新错误类型
|
||||||
|
|
||||||
|
```python
|
||||||
|
class CallDeadlineExceeded(PolyGatewayError):
|
||||||
|
"""调用方设定的整体期限到期;不是网关不可用、也不是源故障。"""
|
||||||
|
def __init__(self, *, scope: str, deadline_s: float) -> None:
|
||||||
|
super().__init__(f"{scope} 调用期限 {deadline_s}s 到期")
|
||||||
|
self.scope, self.deadline_s = scope, deadline_s
|
||||||
|
```
|
||||||
|
|
||||||
|
| 决策 | 取法 | 理由 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 父类 | `PolyGatewayError` **直接**子类 | 不用 `TransientError`: 那是"退避后可重试、计熔断"的源级瞬时故障,下游按它无限外层重试只会把同一份期限再等一遍;不用 `GatewayUnavailableError`: 其消息硬编码"{scope} 网关暂时不可用"(`errors.py:156`),把调用方自选的期限报告成 scope 死亡,正是要避免的"用一个时钟冒充另一个"(BL4) |
|
||||||
|
| `scope` 字段 | **有** | 多 scope 部署时的诊断分组,免得下游从消息串里解析 |
|
||||||
|
| `retry_after_s` 字段 | **无** | 期限到期不含"何时可再试"的信息;给 `0.0` 会按 `errors.py` 既定语义指示下游**立刻重打仍饱和的渠道** |
|
||||||
|
| `SCOPE_REASONS` | **不新增值** | 它不是 `GatewayUnavailableError` 家族成员,与 `reason` 无关 |
|
||||||
|
| 四分类 | 不变 | 它是**调用方策略**的终止信号,不是四分类里的失败;文档须明写 |
|
||||||
|
| 导出 | 进 `__init__.py` 的 `__all__` | 下游要能 `except CallDeadlineExceeded` |
|
||||||
|
|
||||||
|
**醒目**: `except GatewayUnavailableError` / `except AllSourcesExhausted` 的存量代码**接不住**本异常——这是有意设计,且只在显式配置期限后才可能出现。CHANGELOG / wiki / README 必须以此措辞列出。
|
||||||
|
|
||||||
|
## 5. 时钟与失败模式的边界
|
||||||
|
|
||||||
|
### 5.1 两个时钟不混用
|
||||||
|
|
||||||
|
| 时钟 | 用途 | 纪律 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 事件循环时钟(`loop.time()`,`asyncio.timeout` 内部) | 期限的**唯一**计时源 | 只传**相对时长** `deadline_s`;严禁把 `_CallContext._started`(注入 `now`)加上偏移当绝对截止时刻传进去 |
|
||||||
|
| 注入 `now`(`types.py:335`、三个 client) | `CallStats.total_latency_ms`、`StallClock`、退避 | 不读、不改;测试替换它不会影响期限判定,这一点必须在测试里明确 |
|
||||||
|
|
||||||
|
代价写实: 两者不同源,故 `total_latency_ms` 与 `deadline_s` 之间存在微小偏差(注入钟被伪造时可任意大)。这是**有意**的——统一它们要么强迫调用方注入 loop 钟,要么自建定时器,两者都比这点偏差贵。
|
||||||
|
|
||||||
|
### 5.2 四种"到期周边形态"必须分开
|
||||||
|
|
||||||
|
| 现象 | 判据 | 结果 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 本层期限到期 | `except TimeoutError` 且 `cm.expired()` | `CallDeadlineExceeded`,终态行记 deadline |
|
||||||
|
| 内层自抛 `TimeoutError`(未到期) | `cm.expired()` 为假 | 原样上抛,不吞不改判(同 `streaming.py:52-58`);探针 D4 |
|
||||||
|
| **到期后清理路径自抛 `TimeoutError`** | `cm.expired()` 为真但异常对象**就是体内那一个**(身份比较) | 原样上抛该 `TimeoutError`,**不改标成 deadline**(探针 D2/D3;只看 `expired()` 的写法在此会误标)。代价: 该异常不是 `PolyGatewayError`,三个边界的 `except PolyGatewayError` 接不住→**无终态行**(与 1.3.5 已有的裸 `TimeoutError` 穿透行为同口径,非本版新增) |
|
||||||
|
| 外部取消 | `asyncio.timeout` 契约:不是本层计时器造成的取消 → `CancelledError` 原样上抛 | 走既有 `except asyncio.CancelledError` 分支,终态行记 `"cancelled"`;**不会**同时出现两条终态行(`claim_terminal()` 去重)。探针 D7/D8 实测: 外部取消**无论先于还是晚于到期**(含清理期到达)都胜出 |
|
||||||
|
| **到期窗口内体内先抛领域异常** | 计时器已触发、取消尚未投递到达时,体内先 `raise AllSourcesExhausted(...)` 等 | **领域异常原样逐层上抛,期限静默让位**;终态行记该领域异常而非 deadline。该窗口在重试循环真实存在(一次尝试刚结束与计时器同刻),复审 B 探针 E2/E8 已实测 |
|
||||||
|
|
||||||
|
故本设计只承诺"到期**通常**得 `CallDeadlineExceeded`",**不承诺 100%**;实现与测试均不得写成无条件断言。擦边竞态(计时器已触发但调用体已成功返回)不会遗留游离取消——这是 `asyncio.timeout` 的公共行为,库不自判形态、也不依赖任何 CPython 私有实现;跨版本保证由受支持 Python 矩阵上的行为测试提供(§10)。
|
||||||
|
|
||||||
|
### 5.3 期限治理的是"等待",不是返回时刻(必须写进 wiki,不得含糊)
|
||||||
|
|
||||||
|
asyncio 是**合作式取消**: 到期只是向任务投递一次取消,真正返回的时刻取决于在途 `await` 何时到达取消点,以及**清理路径**跑多久。已知会在期限之后继续跑的三段:
|
||||||
|
|
||||||
|
| 段 | 位置 | 性质 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| permit 结算与释放 | `retry.py:341` → `admission.py:46-60` | Redis 后端是两次网络往返;不受期限管辖 |
|
||||||
|
| 取消路径的 attempt 遥测行 | `retry.py:320-323` | 一次落库;丢了就丢了本次现场 |
|
||||||
|
| 终态遥测行 | 三个边界的 `emit_terminal_once` | **有意留在期限之外**:诊断行如果自己被期限切掉,期限到期这件事就没有台账 |
|
||||||
|
|
||||||
|
**量级不是毫秒级**: 复审 B 探针 E4 以上述真实形态(取消分支 0.2s 遥测 + `finally` 0.1s 结算)实测:期限 0.05s → 返回时刻 0.351s(**约 7 倍**);本轮独立探针 D6(清理 0.2s)复现同一形态: 0.251s(**5 倍**)。故对外措辞必须是"**返回时刻 = 期限 + 清理耗时**",而不是"最多 N 秒返回";清理耗时取决于 permit/遥测后端,可远超期限本身。§10 以量化断言把越限量变成可观测、可回归的量。
|
||||||
|
|
||||||
|
**反向代价(同等重要)**: 成功之后的旁路 IO **在期限之内**——`middleware/cache.py:199` 的 `_safe_set` 在 `call_next` 返回之后执行,`middleware/telemetry.py:734` 的缓存命中写同理。期限落在这两步 → 一个已完成、**已计费**的 `LLMResponse` 被丢弃,调用方只拿到 `CallDeadlineExceeded`。故 wiki 必须写明"期限到期**不等于**未产出、未计费";本版**不引入 `shield`** 去抢救它。
|
||||||
|
|
||||||
|
取舍是显式的: **宁可超出期限也要留下资源清理与诊断**,而不是引入 `shield`/后台任务去"抢救"(那会把取消语义弄脏,违反"取消可穿透"铁律)。若下游端口实现(自实现 transport / 遥测)在清理里长时间阻塞,期限的超出量就是那段阻塞时长——库不为第三方实现兜底。
|
||||||
|
|
||||||
|
## 6. 与既有算法的关系(明确不改的三件事)
|
||||||
|
|
||||||
|
### 6.1 F1:`Retry-After` 非有限值(纯 bug 修复)
|
||||||
|
|
||||||
|
`_parse_retry_after` 加判据: **解析成功但为无穷**(`inf`/`-inf`/`1e999`,判据 `math.isinf(seconds)`)时显式忽略并记一条 `warning`,按"服务端没给提示"处理,不伪造缺省值。`nan` 维持现状——它被既有的 `seconds > 0` 恒假拦下,**静默 `None` 且不告警**(§2),本版**不给它加告警、不改判据顺序**。
|
||||||
|
|
||||||
|
告警要有主语但不得回显不可信输入,故函数签名改为 `_parse_retry_after(raw: str | None, *, source_name: str) -> float | None`:`source_name` 是**必填 keyword-only 私有参数**(带前导下划线的模块内函数,不属公共面,无需 keyword 默认值兜底),唯一调用处 `_translate_429`(`transports/openai_compat.py:140`)传 `source_name=source.name`。warning 只写源名与判据词(如 `retry_after_not_finite`),**不拼接、不截断、不打印原始头字符串**。其余形态(HTTP-date、空串、负数、不可解析)**维持静默返回 `None`**--HTTP-date 是 RFC 7231 合法形态、空/负是常见噪声,逐次 warning 会在 429 风暴时把真缺陷的信号淡掉。不抛 `RequestRejectedError`:一个坏响应头不该把一次可重试的 429 判死。
|
||||||
|
|
||||||
|
### 6.2 有限大值的 `Retry-After` 不夹上限
|
||||||
|
|
||||||
|
审查 BL1 建议 `min(retry_after, backoff_max_s)`,**本设计不采纳**:夹小的直接后果是提前重打一个明确说了"3600 秒后再来"的饱和渠道,把服务端调度指令改写成库的猜测。有限大值的处置只有一条正路——调用方设期限,由 §4 的硬边界在到期时切断那次 sleep。未配期限即维持 1.3.5 语义(等满 `Retry-After`),这一残余必须在 wiki 明写。
|
||||||
|
|
||||||
|
### 6.3 取消路径结算修复(**已获批,本版第一个原子提交**)
|
||||||
|
|
||||||
|
现状(现读): `retry.py:281 actual = 0` → `:320` 取消分支 → `:341 finally` → `admission.py:46-60 settle_and_release(permit, actual)`;`settle` 算的是 `delta = actual - est`(`backends/memory/limiter.py:48-55`、`backends/redis/limiter.py:136-158`),故 `actual=0` = **把入场预扣的 TPM 整笔退还**。取消发生在 transport 在途时,上游可能已经计费——退款就是把已消耗的额度退回闸里。启用期限后库自己会常规性触发该路径,故先修后启用。
|
||||||
|
|
||||||
|
**精确结算矩阵**(以“这一刻库到底知道什么”为唯一判据;`est` 指 `source.effective_est_tokens()`):
|
||||||
|
|
||||||
|
| # | 取消落点(精确位置) | 库此时知道的事实 | `settle()` 取值 | 本版是否改变 |
|
||||||
|
| --- | --- | --- | --- | --- |
|
||||||
|
| S1 | 准入阶段(`admission.pick`:`try_enter` 异常、开路分支) | **确定未调用 transport** | `0`(全额退) | 否(既有行为即此) |
|
||||||
|
| S2 | 退避 `sleep` / 配额轮询 / 熝断等待 | 本轮未持 permit(上一轮已在 `finally` 结清) | 无 permit 可结 | 否 |
|
||||||
|
| S3 | `_attempt` 内、transport 在途(`retry.py:288` 的 await 未返回) | **端口已开始但用量未知** | `est`(delta==0,保留预扣) | **是**——原为 `0` |
|
||||||
|
| S4 | transport 已返回、`actual` 已算出后的任一 await(记账写回/逐次遥测) | **完整 usage 已知**(measured/estimated) | 保留已算出的真实 `actual`,**不得被取消分支覆盖回 `est`** | 否(现行为已正确,修复不得弄坏) |
|
||||||
|
| S5 | **已处理领域失败**分支内的 await(`record_failure`/`_emit`)中途取消 | 已知失败类型,结算决定已算出 | 该分支的决定值:瞬时 = `est`;`SourceDead` = **`0`**;`RequestRejected`/`ResultInvalid` = `0` | 是——决定移到该分支同步处理之后、紧贴首个 await 之前,故 `SourceDead` 的既有 `0` 在取消下**被保住**而不再退化成 `est` |
|
||||||
|
| S6 | OCR 任何位置(`ocr.py:449 settle_and_release(permit, 0)`) | **OCR 无 token 是事实**,不是“未知” | `0` | 否(事实即 0,不得改成 `est`) |
|
||||||
|
| S7 | embedding 与 chat 同构两处(`embedding.py:392-407` 取消分支) | 同 S3/S4 | 同 S3/S4 | **是**(与 chat 同口径同时改) |
|
||||||
|
| S8 | **未被任何 except 接住**的异常(`RuntimeError`、`KeyError` 等未分类逃逸) | 库对用量一无所知,且**不在本版批准范围** | `0`(与 1.3.5 逐字一致) | 否——本版**不**把"端口开始 = 可能已计费"推广到未分类异常 |
|
||||||
|
|
||||||
|
**实现形态**(实施时不得变形;人类只批准了"取消路径"这一条,故语义扩大**必须**被限制在取消分支内):
|
||||||
|
|
||||||
|
`actual` 初值**保持 `0` 不动**;另设一个**局部阶段变量** `settlement_known: bool = False`(纯函数内局部,不是新公共面、不进任何签名、不进配置)。置位规则只有两条——成功路径拿到用量(真实 usage 或 `usage_source == "unavailable"` 的 `est`)后置 `True`;三个**已处理领域失败**分支在**进入分支后的第一条语句**算出结算值并置 `True`。取消分支只在 `settlement_known` 仍为 `False` 时才赋 `actual = est`。
|
||||||
|
|
||||||
|
```python
|
||||||
|
actual = 0
|
||||||
|
settlement_known = False # 局部阶段变量: 该刻库是否已算出确定结算
|
||||||
|
# 成功: actual = 真实 usage 或 est → settlement_known = True
|
||||||
|
# RequestRejected / ResultInvalid: actual = 0; settlement_known = True(分支首句)
|
||||||
|
# SourceDead / Transient: actual = 0 if dead else est; settlement_known = True(分支首句)
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
if not settlement_known: # 端口已开始、结算未定 → 保守保留预扣
|
||||||
|
actual = source.effective_est_tokens()
|
||||||
|
...
|
||||||
|
```
|
||||||
|
|
||||||
|
三条不变量(同级 `except CancelledError` 接不住其他 `except` 块内的取消,故必须在其首个 await 前先确定结算,而非只靠标志位):① **已确定的值一律不覆写**,包括真实 usage 恰为 `0`(源真返回 0 token 是事实,不是"未知");② "确定结算"的界桩按**当前代码里 failure 记账 await 的前后位置**划定——失败分支的结算决定被前移到 `record_failure` / `_emit` **之前**,故取消无论落在这两个 await 的哪一侧,拿到的都是该失败类型本来的值;③ 未被任何 `except` 接住的异常不经上述任一分支,`actual` 保持 `0` 原样传播(S8)。
|
||||||
|
|
||||||
|
**明确收窄**: 本版**不再**把 S5 泛化成"一切失败按 `est` 结算"。1.3.5 的 `SourceDead` 退全款(`0`)是**有意**的语义(源已判死,不该继续占额度),取消恰好落在它之后时必须保留那个 `0`;"端口开始 = 可能已计费"只是**取消且结算未定**这一格的兜底取值,不是全局通则。
|
||||||
|
|
||||||
|
这样只有取消路径的取值发生改变,其余四条既有路径与未分类异常路径逐字不变(验收矩阵 §10 "结算"批次逐条钉,并新增一条**防越界回归**)。
|
||||||
|
|
||||||
|
**保守结算的诚实边界**(不得写成“修对了”): `await transport.complete(...)` 返回前被取消,只能证明**端口协程已被进入**,不能证明 HTTP 字节已发出、更不能证明上游已计费。故 S3 是一个**保守选择**(宁可多扣不可凭空退款),不是事实性计量;它与既有的“瞬时失败按 `est` 结算”(`retry.py:337`)同一口径、同一理由。库**不新增任何 wire 事件协议**(如“transport 上报字节已发出”)来缩小这个不确定区——那是新公共端口面,且 httpx 层面也给不出可靠信号;不确定性写进文档而不是藏起来。
|
||||||
|
|
||||||
|
**不改且不被本版解决的相邻缺口**(列出以免被读成已修): ① `ResultInvalidError` 路径(如 `embedding.py:350-355` 维度不符、`transports/openai_compat.py:280-308` 响应形态异常)——响应真实返回过(已计费)但仍按 `0` 退全款;② `RequestRejectedError` 同理。两者与取消无关,属另一族记账语义变更,**未获批准即不动**,另行立 issue。
|
||||||
|
|
||||||
|
**共享状态不变**: `permit.release()`、`pacer.leave()`、`breaker.release_probe()`(探针归还)、`mark_progress` 全部保持原样且仍在 `finally`;本修复**只改 `settle()` 的入参取值**,不动限流 Lua、不动端口签名、不动幂等语义。
|
||||||
|
|
||||||
|
## 7. 遥测(无新增列,无新增 DDL)
|
||||||
|
|
||||||
|
| 行 | 到期时的取值 |
|
||||||
|
| --- | --- |
|
||||||
|
| attempt 行 | 被取消的那次仍记 `error="cancelled"`(`retry.py:320-323`)——它描述的是**那次尝试**的真实结局,保留 |
|
||||||
|
| 终态行 | 经既有 `emit_terminal_once` 写出;到期**通常** `error_type='CallDeadlineExceeded'`(例外见 §5.2 第 4 行:体内先抛领域异常时记该异常)、`error` 为其消息串(`telemetry.py:274-305` 既有取值逻辑,零改动),`logical_call_id`/`attempts`/`total_latency_ms` 照旧 |
|
||||||
|
| 计数 | 每逻辑调用至多一条终态行,由 `claim_terminal()` 保证;**不会**同时出现 cancelled 与 deadline 两条 |
|
||||||
|
|
||||||
|
## 8. 变更点清单(反 gold-plating)
|
||||||
|
|
||||||
|
| 类别 | 内容 |
|
||||||
|
| --- | --- |
|
||||||
|
| 新增文件 | `src/polygateway/deadline.py`(1 个校验函数 + 1 个包裹函数) |
|
||||||
|
| 改动文件 | `errors.py`(1 个类)、`__init__.py`(1 个导出)、`config.py`(1 个键 + 1 个 loader + 1 个守卫 + 1 个字段,守卫直调 `deadline.ensure_call_deadline`)、`client.py`/`embedding.py`/`ocr.py`(各 1 处包裹 + 构造参数 + 公开方法参数 + 入口校验 + `from_settings` 透传)、`middleware/retry.py` 与 `embedding.py` 的 `_attempt`(§6.3 结算矩阵:只改 `actual` 取值 + 1 个**局部**阶段变量 `settlement_known`)、`transports/openai_compat.py`(F1 一行判据 + warning + `source_name` 必填私有 kw 及其唯一调用处)、`pyproject.toml`(import-linter layers 新增 `polygateway.deadline` 一层) |
|
||||||
|
| 直接复用 | `_CallContext`、`emit_terminal_once`、`claim_terminal`、`asyncio.timeout` + `cm.expired()` 范式、`settle_and_release` 单一出口、既有 FakeClock/假 transport 测试设施、限流契约套件(Lua 不改) |
|
||||||
|
| **明确不做** | 不改 `backoff_delay`/`SourceAdmission._nap`/`StallClock`/429 分账/限流 Lua/端口签名;不改 `ResultInvalid`/`RequestRejected` 的结算口径;不加遥测列;不加 `shield`/后台任务;不加 429 次数键;不做 #24 |
|
||||||
|
|
||||||
|
## 9. 集中人类批准项(2026-09-10 全数获批)
|
||||||
|
|
||||||
|
| # | 决策 | 结果 | 不采纳的代价 |
|
||||||
|
| --- | --- | --- | --- |
|
||||||
|
| H1 | 采纳 §3 方案 A(单一 `asyncio.timeout` 硬边界),缺省 `None` | **已批** | B/C 只能给轮次粒度或多键相加的"伪期限" |
|
||||||
|
| H2 | `CallDeadlineExceeded` 为 `PolyGatewayError` 直接子类、无 `retry_after_s`、不进 `SCOPE_REASONS` | **已批** | 挂进 `GatewayUnavailableError` 会把调用方的期限报告成网关不可用 |
|
||||||
|
| H3 | 取消结算修复的位置 | **已批 (a′)**: 不再另立前置 issue,改为 **1.3.6 内的第一个原子提交**,口径按 §6.3 矩阵(S1-S8,含 S8 未分类异常维持 `0` 的收窄)逐格定死 | 选 (b) = 把已知记账缺口变成常规路径,正是本项目反复吃过的亏 |
|
||||||
|
| H4 | per-call 覆盖形态: §4.3 的"`None` 继承、无单次关闭" | **已批** | 三态哨兵扩大公共面;完全不给 per-call 则长短调用必须装两个 client |
|
||||||
|
| H5 | §6.2 不夹有限大 `Retry-After`(与审查 BL1 建议相反) | **已批**,残余写进 wiki | 夹小即提前重打饱和渠道 |
|
||||||
|
| H6 | §5.3/§11 的对外承诺形态: 期限治理**等待**,返回时刻 = 期限 + 清理耗时(实测可达 5-7 倍),且到期可能丢弃已计费成功 | **已批**,不引入 `shield` | 写成"最多 N 秒返回"会让下游上层超时被整片击穿 |
|
||||||
|
| H7 | issue #22 关闭判据 = "调用方**能配置**上限";#24 保持 open、**本版不实现** | **已批** | 见 §11 |
|
||||||
|
|
||||||
|
实施边界不得再扩: 本表之外的任何公共面变化(新配置键、新端口方法、新遥测列、其它记账口径变更)均属**未批准**,需停下来报。
|
||||||
|
|
||||||
|
## 10. 离线验收矩阵(实施时须先失败后通过;全部不触网、不付费)
|
||||||
|
|
||||||
|
| 批次 | 断言 |
|
||||||
|
| --- | --- |
|
||||||
|
| 未启用回归 | `call_deadline_s=None` 时三条链路行为逐字不变:现有 `tests/unit/test_retry.py`、`test_backpressure.py`、`test_embedding.py`、`test_ocr_client.py`、`test_client.py` 全绿(不改一行断言) |
|
||||||
|
| 期限命中 | 假 transport 真实 `await asyncio.sleep(0.3)`、`call_deadline_s=0.05` → 抛 `CallDeadlineExceeded`,`scope`/`deadline_s` 正确;计时范式沿用 `tests/unit/test_streaming.py` 的真实 loop 时钟 + 4-10× 余量(已验证稳定,不标 slow) |
|
||||||
|
| 覆盖面 | 分别令 ①退避 sleep(注入真 `asyncio.sleep`)②准入排队(配额满轮询)③结构化重问 ④embedding 多批 各自超期 → 均抛 `CallDeadlineExceeded`;embedding 断言 N 批**共享一份**期限(总时长不随批数放大) |
|
||||||
|
| 形态区分 | ①内层自抛 `TimeoutError`(假 transport 直抛)且未到期 → 原样上抛,不变成 deadline;②外部 `task.cancel()` → 仍抛 `CancelledError`,终态行 `error="cancelled"`;③擦边成功(transport 耗时略小于期限)→ 正常返回,无游离取消;④**到期窗口内体内先抛领域异常**(自旋构造确定性窗口)→ 上抛该领域异常、终态行记它,**不断言必为 deadline**;⑤**到期后清理自抛 `TimeoutError`**(假端口在取消分支里 `raise TimeoutError`)→ 原样上抛该 `TimeoutError`,**断言不是 `CallDeadlineExceeded`**(钉住身份比较机制;只看 `cm.expired()` 的写法在此变红) |
|
||||||
|
| 结算(§6.3 矩阵逐格) | S1 准入取消 → `tpm_used` 回到 0;S3 transport 在途取消 → `tpm_used == est`(**先失败后通过**的核心红绿);S4 取消落在 usage 已知之后 → `tpm_used == 真实 usage`(不被 `est` 覆盖);S6 OCR 取消 → `tpm_used == 0`;S5 取消落在 `SourceDead` 的 `record_failure` await 中途 → `tpm_used == 0`(**不得**变成 `est`);**S8 防越界回归**: 假 transport 抛 `RuntimeError`(未分类)→ `tpm_used == 0` 且异常原样上抛;成功/`SourceDead`/`RequestRejected`/`ResultInvalid` 四条既有路径结算值逐字不变;真实 usage 恰为 0 的成功 → `tpm_used == 0`;permit `release`/`pacer.leave`/`release_probe` 调用次数不变 |
|
||||||
|
| 终态遥测 | 到期恰好一条 `event_kind='terminal_failure'` 行,`error_type='CallDeadlineExceeded'`;被取消的 attempt 行仍为 `cancelled`;两者 `logical_call_id` 一致;列数不变 |
|
||||||
|
| 清理不可越过(含量化) | 到期时 permit 的 `settle`/`release` 与 pacer `leave` 仍被调用(假 permit 记账断言);`CancelledError` 未被吞;**量化断言**: 注入已知 sleep 的假 permit/假 emitter → 返回时刻 ≈ 期限 + 已知清理时长(可远大于期限本身) |
|
||||||
|
| 成功旁路被切 | 假缓存后端 `set` 慢于期限 → 抛 `CallDeadlineExceeded` 且断言上游 transport **已成功调用一次**(已计费成功被丢弃,§5.3),把该行为钉住而非留作偶然 |
|
||||||
|
| 注入钟无关 | 伪造注入 `now`(跳变 10^6 秒)不触发期限;反之期限触发时 `total_latency_ms` 仍来自注入钟 |
|
||||||
|
| 配置(全装配路径) | 键未设 → `None`;`0`/负/`nan`/`inf`/非数/`True`(bool 不得当 1 秒)→ `ValueError` 且消息含来源——**四条路各测一遍**: env、`GatewaySettings` 直接构造、`dataclasses.replace`、**三个 client `__init__` 直传**;per-call 非法值在构造协程前抛错且**无 "coroutine was never awaited" 警告**;`call_deadline_s < timeout_s` **不报错**(允许) |
|
||||||
|
| F1 | `Retry-After` 为 `inf`/`-inf`/`1e999` → 按无提示处理且各记一条 warning(**断言日志含源名、不含原始字符串**);`nan`/空/负/HTTP-date → 静默 `None` 且**无 warning**;`_parse_retry_after` 的 `source_name` 为必填 kw(漏传即 `TypeError`);有限正数仍取大;`insufficient_quota` 仍归 `SourceDead` |
|
||||||
|
|
||||||
|
## 11. issue 闭环判据
|
||||||
|
|
||||||
|
| issue | 判据 | 措辞纪律 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| #22 | **可关闭**: 调用方通过 `{SCOPE}__CALL_DEADLINE_S` 或 per-call 参数即可给一次调用设上限 | 措辞必须是"**期限治理的是等待,不是返回时刻**;返回时刻 = 期限 + 清理耗时(取决于 permit/遥测后端,实测可达数倍)",不得写成"最多等 N 秒"的硬保证;同时写明"到期不等于未产出、未计费"(§5.3)与"**不配置就保持 1.3.5 旧语义**"(纯 429 序列仍可能长时间等待、有限大 `Retry-After` 仍照睡)——三句均需出现在 CHANGELOG、wiki 与 issue 关闭说明的显要位置 |
|
||||||
|
| #24 | **保持 open、本版不实现**(未获批准),另行设计 | 期限只让长尾**更早失败**,与"更快成功"是两件事;任何文档不得把前者写成后者,也不得因本版而声称 #24 缓解 |
|
||||||
|
|
||||||
|
## 12. 待补证据与残余风险(诚实标注)
|
||||||
|
|
||||||
|
| 项 | 状态 |
|
||||||
|
| --- | --- |
|
||||||
|
| 取消是否让上游停止生成与计费 | 无一手证据 → §6.3 的 S3 只能是**保守选择**而非事实计量;不做任何"取消即省钱"或"已修准"的表述 |
|
||||||
|
| 保守结算的反向偏差 | S3 在“transport 确实未发出字节”时会**多扣** `est`(直到窗口滞后自然过期);与“凭空退款击穿网关”相比这是有意选定的方向(降级方向铁律),但必须写进 CHANGELOG |
|
||||||
|
| 未分类异常的结算 | `RuntimeError` 等未被四分类接住的异常仍按 `0` 退全款(S8)——人类只批准了取消路径,本版**不**将保守口径扩到它们;它们理论上同样可能发生在 transport 在途之后,属已知残留,需时另立 issue |
|
||||||
|
| `ResultInvalid`/`RequestRejected` 的退全款 | 本版**不改**(§6.3 末段);它们与取消无关,属未批准的另一族记账语义变更 |
|
||||||
|
| 跨 Python 版本的取消/超时语义 | 本设计只依赖 `asyncio.timeout` 的**公共行为**与异常对象身份比较,不引 CPython 私有实现作保证;保证由 §10"形态区分"在受支持 Python 矩阵上的测试提供(本轮探针仅覆盖 3.12.13) |
|
||||||
|
| 清理自抛 `TimeoutError` 时无终态行 | 该异常不是 `PolyGatewayError`,三个边界的 `except` 接不住(§5.2);与 1.3.5 已有的裸 `TimeoutError` 穿透同口径,本版不扩大也不修补 |
|
||||||
|
| #22 现场 46.7s/20min 数字 | 未复跑;本设计不依赖其数值,只依赖路径成立性(已由源码证明) |
|
||||||
|
| 第三方端口实现的清理耗时 | 不可控,直接构成期限超出量(§5.3) |
|
||||||
|
| 期限与 `stall_window_s` 的联合调参建议 | 本版不给推荐值:两者治理对象不同(调用方意志 vs scope 活性),给一个"经验公式"就是把两个预算再次绑死 |
|
||||||
@@ -0,0 +1,109 @@
|
|||||||
|
---
|
||||||
|
type: finding
|
||||||
|
node_id: finding:2026-09-10-136-call-deadline-validation
|
||||||
|
title: "1.3.6 可选调用期限与取消结算验证记录"
|
||||||
|
date: 2026-09-10
|
||||||
|
---
|
||||||
|
|
||||||
|
# 1.3.6 可选调用期限与取消结算验证记录
|
||||||
|
|
||||||
|
> 范围:分支 `feature/1.3.6-call-budgets` 上的 T1–T3 三个行为提交(`1ff83bb` / `9474c76` / `da77b12`)与 T4 文档提交。设计 `research-wiki/designs/2026-09-09-136-call-budgets-design.md`(人类已批准),计划 `research-wiki/plans/2026-09-10-136-call-deadline.md`。
|
||||||
|
> 本文件是**证据索引**:原始输出在 `tests/outputs/136/`(按纪律**不提交**),此处只记路径、命令、退出码与结论。
|
||||||
|
> 环境:conda 环境 `PolyGateway`,**Python 3.12.13**(`conda run -n PolyGateway python -V` 实测)。所有 pytest/lint 命令均**不接管道**,退出码直取。
|
||||||
|
> 版本号未 bump、未 tag、未发布——发布清单(CLAUDE.md §4.4.1)不在本轮范围。
|
||||||
|
|
||||||
|
## 1. 红绿证据索引
|
||||||
|
|
||||||
|
| 任务 | 阶段 | 证据文件 | 结果 |
|
||||||
|
| --- | --- | --- | --- |
|
||||||
|
| T1 取消结算 | 红/绿 | **未落盘**(见 §1.1 诚实说明) | 见 §1.1 |
|
||||||
|
| T1 真实 Redis | 绿(本轮在 HEAD `da77b12` 上复跑) | `tests/outputs/136/t4/redis_cross_connection.txt` | `8 passed`,`exit=0` |
|
||||||
|
| T2 期限 | 红(值域+形态) | `tests/outputs/136/t2/red_deadline.txt` | 收集期 `1 error`(`deadline.py` 缺席),`exit=2` |
|
||||||
|
| T2 期限 | 绿(`test_deadline.py`) | `tests/outputs/136/t2/green_deadline.txt` | `22 passed`,`exit=0` |
|
||||||
|
| T2 接线中途 | 绿 | `tests/outputs/136/t2/unit_contracts_midway.txt`、`unit_contracts_after_wiring.txt` | 各 `1560 passed, 17 skipped` |
|
||||||
|
| T2 入口冒烟 | 绿 | `tests/outputs/136/t2/entry_smoke.txt` | 四个方法签名含 `call_deadline_s`;到期异常 `has retry_after_s: False`,`exit=0` |
|
||||||
|
| T2 回归门 | 绿 | `tests/outputs/136/t2/final_unit_contracts.txt` | `1572 passed, 17 skipped`,`pytest_exit=0` |
|
||||||
|
| T2 lint | 红→绿 | `tests/outputs/136/t2/lint.txt`(`Found 3 errors`,`lint_exit=2`)→ `lint_final.txt`(`All checks passed!` + `Contracts: 1 kept, 0 broken.`,`lint_exit=0`) | 修后绿 |
|
||||||
|
| T2c 补测 | 红 pass1 | `tests/outputs/136/t2c/red_pass1_import_absent.txt` | `3 errors in 0.29s`(三个模块收集期 ImportError) |
|
||||||
|
| T2c 补测 | 红 pass2 | `tests/outputs/136/t2c/red_pass2_real_reasons.txt` + `red_method_and_reasons.txt` | `32 failed`;分布见 §1.2 |
|
||||||
|
| T2c 补测 | 绿 | `green_unit_after_hardening.txt`(`1530 passed`)、`green_client_recheck.txt`(`110 passed`,`EXIT=0`)、`green_recheck_client_embedding.txt`(`165 passed`)、`green_unit_contracts.txt`(`1592 passed, 17 skipped`)、`final_unit_contracts.txt`(`1606 passed, 17 skipped`) | 全绿 |
|
||||||
|
| T2c lint | 绿 | `tests/outputs/136/t2c/lint.txt` | `All checks passed!` + `Contracts: 1 kept, 0 broken.` |
|
||||||
|
| T3 `Retry-After` | 红 | `tests/outputs/136/t3/red.txt` | `6 failed, 8 passed, 129 deselected` |
|
||||||
|
| T3 `Retry-After` | 绿 | `green_file.txt`(`143 passed`)、`green_unit_contracts.txt`(`1606 passed, 17 skipped`) | 全绿 |
|
||||||
|
| T3 lint | 绿 | `tests/outputs/136/t3/lint.txt` | `All checks passed!` + `Contracts: 1 kept, 0 broken.` |
|
||||||
|
|
||||||
|
`tests/outputs/136/t2/lsp_noise_refutation.txt` 与 `t2c/lsp_noise_refutation.txt` 记录编辑器 LSP 报的 import/属性告警属环境噪声(`pydantic` 在 conda 环境可解析),不是代码缺陷。
|
||||||
|
|
||||||
|
### 1.1 T1 的证据形态(诚实说明)
|
||||||
|
|
||||||
|
T1(`1ff83bb`)的**先红后通过证据产生于当时的会话工具输出,未落盘为 `tests/outputs/136/t1/` 文件**。本文件不追认那次输出,只登记两项**当下可复核**的替代证据:
|
||||||
|
|
||||||
|
| 替代证据 | 内容 |
|
||||||
|
| --- | --- |
|
||||||
|
| 提交 `1ff83bb` 的 diff | 三个源文件 + 四个测试文件共 213 插入;测试侧含 S3/S7(取消结算按 `est`)、S5-dead(仍 `0`)、S8(`RuntimeError` 逃逸仍 `0`)、真实 usage 恰为 0 的成功仍 `0` 四组断言 |
|
||||||
|
| 本轮在 HEAD `da77b12` 上重跑真实 Redis | `pytest tests/integration/test_redis_cross_connection.py -q` → `8 passed`,`exit=0`(`tests/outputs/136/t4/redis_cross_connection.txt`) |
|
||||||
|
|
||||||
|
结论口径:**T1 的“红”只有会话内证据、无归档文件**;T1 的“绿”在当前 HEAD 上已被真实 Redis 复现证实。
|
||||||
|
|
||||||
|
### 1.2 T2c 两趟红证据的方法说明(诚实说明)
|
||||||
|
|
||||||
|
T2c 的红证据是在**基线 `1ff83bb`(T2 之前)**的 `git worktree --detach` 检出上取的,用 `PYTHONPATH=<worktree>/src` 覆盖 editable `.pth`(已实测 `polygateway` 加载自 worktree 且 `deadline.py` 缺席):
|
||||||
|
|
||||||
|
| 趟次 | 做法 | 结果 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| pass1 | 用例原样跑 | 三个测试模块**收集期** ImportError(`CallDeadlineExceeded` 不存在)→ `3 errors`。只证明符号缺席,**没有执行到函数体** |
|
||||||
|
| pass2 | 仅把缺失的**导入符号**替换成**本地占位异常类**(shim),让函数体真正跑起来 | `32 failed`:**27 条“参数/属性不存在”**(`chat()` 9、`embed()` 4、`GatewaySettings.__init__()` 3、`GatewaySettings.call_deadline_s` 属性 3、`recognize_text()` 2、`GatewayClient.__init__()` 2、`parse_layout()`/`OcrClient.__init__()`/`EmbeddingClient.__init__()`/`GatewayClient._call_deadline_s` 各 1)+ **5 条 `Failed: DID NOT RAISE ValueError`** |
|
||||||
|
|
||||||
|
**该 shim 是一次性本地脚手架,未提交、不在任何分支上**;它只替换导入符号,不改被测源码。故 pass2 的红**是针对预 T2 源码的真实失败原因分布**,而非构造错误——但读者需知这份红**无法从仓库检出复现**,只能从上表与 `red_method_and_reasons.txt` 复核。
|
||||||
|
|
||||||
|
## 2. 命令与退出码
|
||||||
|
|
||||||
|
| 命令(前缀均为 `conda run -n PolyGateway python -m`,lint 为 `make lint`) | 何时 | 结果 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| `pytest tests/unit/test_deadline.py -q` | T2 红 | `exit=2`(collection error,符合预期) |
|
||||||
|
| `pytest tests/unit/test_deadline.py -q` | T2 绿 | `22 passed`,`exit=0` |
|
||||||
|
| `pytest tests/unit tests/contracts -q` | T2 门 | `1572 passed, 17 skipped`,`exit=0` |
|
||||||
|
| `pytest tests/unit/test_client.py tests/unit/test_embedding.py tests/unit/test_ocr_client.py tests/unit/test_config.py -q -rf -k "…deadline…"` | T2c 红 | `32 failed`(基线 worktree,见 §1.2) |
|
||||||
|
| `pytest tests/unit tests/contracts -q` | T2c 门 | `1606 passed, 17 skipped`,`exit=0` |
|
||||||
|
| `pytest tests/unit/test_openai_compat.py -q -rf -k RetryAfterNonFinite` | T3 红 | `6 failed, 8 passed, 129 deselected` |
|
||||||
|
| `pytest tests/unit/test_openai_compat.py -q` | T3 绿 | `143 passed` |
|
||||||
|
| `pytest tests/unit tests/contracts -q` | T3 门 | `1606 passed, 17 skipped`,`exit=0` |
|
||||||
|
| `pytest tests/integration/test_redis_cross_connection.py -q` | T1/T4 复跑 | `8 passed`,`exit=0` |
|
||||||
|
| `make lint`(ruff + import-linter) | T2/T2c/T3 收尾 | `All checks passed!`;`Contracts: 1 kept, 0 broken.`(新 `polygateway.deadline` 层在内) |
|
||||||
|
|
||||||
|
T4(本次文档提交)**不含行为变更**,故未跑测试;仅新增上表最后一行的真实 Redis 复跑作为 §1.1 的替代证据。
|
||||||
|
|
||||||
|
## 3. 豁免索引(哪些门没跑,为什么,谁来兜)
|
||||||
|
|
||||||
|
| 未执行项 | 原因 | 兜底责任 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| `pytest -m slow`(真实网关 e2e、Redis 时间语义变体) | 成败取决于外部服务当下状态,默认被 `addopts = "-m 'not slow'"` 排除;计划 §5 明确本轮不跑 | **发布清单(CLAUDE.md §4.4.1)第 4 步**,合并 main 后统一执行 |
|
||||||
|
| `tests/e2e/` 四个文件 | 同上,本版零付费调用 | 同上 |
|
||||||
|
| 真实网关的期限行为实测 | 期限用例用真实事件循环时钟+假 transport 构造,余量 4–10 倍,不依赖网关 | 发布清单第 4 步的 e2e 顺带覆盖;**本轮无真实网关证据** |
|
||||||
|
| 模型能力矩阵复验 | 本版未触碰推理/能力表 | 不适用 |
|
||||||
|
| 跨 Python 版本验证 | 见 §4 残余三 | 未兜底,登记为残余 |
|
||||||
|
|
||||||
|
真实 Redis **不在豁免之列**:T1 已证、本轮在 HEAD 上复跑(`8 passed`),未以 memory 后端冒充。
|
||||||
|
|
||||||
|
## 4. 残余风险(三条,逐条复述设计 §12)
|
||||||
|
|
||||||
|
| # | 残余 | 诚实口径 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 1 | 取消结算按 `est` 保留预扣 | 这是**保守选择,不是“上游已计费”的证明**。库无法知道端口已开始的那次调用是否真的产生了计费用量;方向定为宁多扣不空退(多扣只损失本窗口一点额度,空退会让已计费用量绕过闸门)。真实 usage 已知(含恰为 0)与已判 `SourceDead` 的 `0` 不被覆写;未分类异常逃逸仍按 `0`,属**已知残留,本版不动** |
|
||||||
|
| 2 | 清理期自抛 `TimeoutError` 时**无终态遥测行** | 与 1.3.5 的裸 `TimeoutError` 穿透**同一口径**(这条路径一直存在、一直没有终态行),本版没有让它变坏;**但期限把这条路径常态化了**——启用期限后触发清理的频率上升,其可达性随之上升。`with_call_deadline` 的局部变量身份比较保证这种 `TimeoutError` **不会**被误标成 `CallDeadlineExceeded`(`test_deadline.py` 有断言钉住) |
|
||||||
|
| 3 | 跨 Python 版本仅 3.12.13 有实证 | `asyncio.timeout` 的 `cm.expired()`/`uncancel()` 行为与清理期异常传播是探针在 **3.12.13 单一版本**上实测的;3.13+的行为未验证。库声明 3.12+,故这是**真实的验证缺口**,不是理论担忧 |
|
||||||
|
|
||||||
|
## 5. 发布清单第 4 步结果(2026-09-10,main `b0ab39e`)
|
||||||
|
|
||||||
|
| 门 | 结果 | 证据 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| `make lint`(合并后 main) | 通过(ruff + import-linter 1 kept 0 broken) | 会话内输出 |
|
||||||
|
| 全套件(unit+contracts+integration) | **1683 passed, 23 skipped, exit=0** | `tests/outputs/136/release/full-gate.log` + `.exit` |
|
||||||
|
| slow 交集子集(Redis 时间语义变体 + 真实网关冒烟) | **22 passed, exit=0**(21 分钟) | `tests/outputs/136/release/slow-scoped.log` + `.exit` |
|
||||||
|
| `test_thinking_live.py` 全模型能力矩阵 | **未跑——按 2026-09-10 人类批准的新规则豁免**:本版 diff 零触碰 `thinking.py`/能力注册表/相关 e2e 设施,复用 1.3.4/1.3.5 已登记矩阵证据;当晚多渠道额度耗尽,全矩阵只会产出超时链噪声。规则变更已写入 CLAUDE.md §4.4.1 第 4 步与 §4.6(提交 `b0ab39e`);被中途终止的全量尝试日志留存于 `slow-gate.log` 备查 | 本表 |
|
||||||
|
|
||||||
|
## 6. 未在本轮做的事
|
||||||
|
|
||||||
|
- 不实现 issue #24(长尾对冲):未获批准,代码与文档均无 hedge 机制,本版**不缓解 #24**。
|
||||||
|
- 不修 `ResultInvalid`/`RequestRejected` 已计费坏结果仍退全款:属另一族记账语义,未批准,登记待立 issue。
|
||||||
|
- 不 bump 版本号、不打 tag、不构建、不上传 registry、不同步 wiki——全部留给发布清单。
|
||||||
@@ -0,0 +1,270 @@
|
|||||||
|
---
|
||||||
|
type: plan
|
||||||
|
node_id: plan:2026-09-10-136-call-deadline
|
||||||
|
title: "1.3.6 可选调用期限与取消结算修复实施计划"
|
||||||
|
date: 2026-09-10
|
||||||
|
---
|
||||||
|
|
||||||
|
# 1.3.6 可选调用期限与取消结算修复实施计划
|
||||||
|
|
||||||
|
> 设计:`research-wiki/designs/2026-09-09-136-call-budgets-design.md`,**人类于 2026-09-10 正式批准**(§9 七项批准项全数获批,H3 取 (a′):取消结算修复并入本版而非另立前置 issue)。
|
||||||
|
> 计划审核门:Claude 自审 + 独立模型审查;plan 无人类门,审毕直接执行。
|
||||||
|
> 目标:① 修好取消路径的 TPM 结算(§6.3 矩阵 S1–S7);② 给一次逻辑调用一条**可选**墙钟硬边界(issue #22),缺省 `None` 时行为逐字等于 1.3.5;③ 修 `Retry-After` 非有限值防御缺口。
|
||||||
|
> 方案:设计 §3 方案 A——三个公开边界各一次 `asyncio.timeout`,配**局部变量身份比较**判据(探针实证:只看 `cm.expired()` 会把清理期自抛的 `TimeoutError` 误标成 deadline)。
|
||||||
|
> 技术:Python 3.12+、asyncio、frozen dataclass、pytest + FakeClock + 真实 `asyncio.Event`、真实实验室 Redis、ruff、import-linter。
|
||||||
|
> 基线 HEAD:`d2455e8`(分支 `feature/1.3.6-call-budgets`;工作区仅 `CLAUDE.md` 既有 markdown 差异、未跟踪 `.pi/` 与本轮两份文档,前两者一律不动、不暂存)。
|
||||||
|
|
||||||
|
**不实现 issue #24(长尾对冲)**:未获批准,任何提交、测试与文档均不得出现 hedge/对冲机制,也不得声称本版缓解 #24。
|
||||||
|
|
||||||
|
## 1. 边界、授权与执行纪律
|
||||||
|
|
||||||
|
| 项目 | 固定边界 |
|
||||||
|
| --- | --- |
|
||||||
|
| 唯一 writer | 一工作区一 writer;父会话负责前台委派与审核派发。1.3.X 合并/发布授权沿用;跨到 1.4、新公共面变化或验证豁免须停下确认 |
|
||||||
|
| 公共面 | 只做设计 §9 已批准四项:新错误类 `CallDeadlineExceeded`、新配置键 `{SCOPE}__CALL_DEADLINE_S`、三个 client 构造参数 + 四个公开方法 keyword-only 参数、取消路径结算口径。**不新增其它键/端口方法/遥测列** |
|
||||||
|
| 记账边界 | 只改**取消路径**的 `settle()` 入参取值(设计 §6.3 S3/S5/S7);成功、`RequestRejected`、`ResultInvalid`、`SourceDead` 四条既有路径与**未分类异常逃逸路径(S8,仍 `0`)**的结算值逐字不变;实现只能用**函数内局部阶段变量**,不得新增公开参数;限流 Lua、`Permit` 端口签名、幂等语义一律不动 |
|
||||||
|
| 依赖铁律 | 新模块 `deadline.py` 只 import stdlib + `errors.py`;`middleware/` 仍只依赖端口与内核;import-linter 契约新增一层执法 |
|
||||||
|
| 取消 | `CancelledError` 永不吞没;不引入 `shield`、不开后台任务;清理仍在 `finally`,允许超出期限 |
|
||||||
|
| 降级方向 | 限流/熔断后端仍 fail-closed;遥测/缓存仍 warning 降级;非法期限值 → **当场 `ValueError`**(装配错误不属降级面) |
|
||||||
|
| 证据与秘密 | 不打印 `.env`、token、Authorization;不提交 `.pi/`、`tests/outputs/`;命令输出只记路径、状态与退出码 |
|
||||||
|
| 证据复用 | 复用既有 FakeClock / 假 transport / `settle_and_release` 出口 / 限流契约套件 / 真实 Redis 跨连接用例;**不重跑模型能力矩阵**,本版零付费调用 |
|
||||||
|
|
||||||
|
Skill 纪律:T0 已执行 `writing-plans`;T1–T3 行为变更执行 `test-driven-development`(先失败后通过的证据须落在本会话工具输出里);每次提交执行 `commit`(英文祈使标题、无 AI 签名、显式路径暂存);T4 前执行 `requesting-code-review` 与 `verification-before-completion`;异常先 `systematic-debugging` 定根因。
|
||||||
|
|
||||||
|
## 2. 文件职责与不变接缝
|
||||||
|
|
||||||
|
| 动作 | 精确路径 | 职责 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| 新建 | `src/polygateway/deadline.py` | `ensure_call_deadline()` 值域校验 + `with_call_deadline()` 单一硬边界(§3.1) |
|
||||||
|
| 修改 | `src/polygateway/errors.py` | 追加 `CallDeadlineExceeded(PolyGatewayError)`(§3.2);`SCOPE_REASONS`/四分类**不动** |
|
||||||
|
| 修改 | `src/polygateway/__init__.py` | `from polygateway.errors import ... CallDeadlineExceeded`;`__all__` 插在 `"CallStats"` 之后、`"CircuitOpenError"` 之前(现读 `:61-62`;该列表并非全字母序,头部 `DEFAULT_PROFILES`/`EFFORT_ORDER`/`Effort` 是既有例外,**不得顺手重排**) |
|
||||||
|
| 修改 | `src/polygateway/config.py` | `GatewaySettings` 末尾追加 `call_deadline_s: float \| None = None`;`_load_call_deadline()`;`_validate_call_deadline()` 进 `__post_init__`(§3.3) |
|
||||||
|
| 修改 | `src/polygateway/client.py` | `__init__` 追加 `call_deadline_s`(入口即校);`chat()` 追加 per-call 参数;`:398` 包裹;`from_settings` 透传 |
|
||||||
|
| 修改 | `src/polygateway/embedding.py` | 同上三处(`:209` 包裹整次 `_embed_all`);`_attempt` 结算矩阵(§3.5) |
|
||||||
|
| 修改 | `src/polygateway/ocr.py` | `__init__`/两个公开方法/`_call` 两级透传;`:274` 包裹 `_run`;`:449 settle_and_release(permit, 0)` **保持 0** |
|
||||||
|
| 修改 | `src/polygateway/middleware/retry.py` | `_attempt` 结算矩阵(§3.5:`actual` 初值仍 `0` + 局部 `settlement_known`);`__call__` 循环、`backoff_delay`、`StallClock` 一字不动 |
|
||||||
|
| 修改 | `src/polygateway/transports/openai_compat.py` | `_parse_retry_after`(`:111-119`)增 `math.isinf` 判据 + 一条 warning + 必填私有 kw `source_name`;同步唯一调用处 `_translate_429`(`:140`)(§3.6) |
|
||||||
|
| 修改 | `pyproject.toml` | import-linter layers 在 `"polygateway.thinking"` 与 `"polygateway.providers : polygateway.sources"` 之间插入 `"polygateway.deadline"` 一行 |
|
||||||
|
| 新建 | `tests/unit/test_deadline.py` | `deadline.py` 的值域与五种形态区分(§5 批次 A/B) |
|
||||||
|
| 修改 | `tests/unit/test_retry.py` | 取消结算红绿(S1/S3/S4/S5)+ `FakeTransport` 加 `entered` Event |
|
||||||
|
| 修改 | `tests/unit/test_embedding.py` | 取消结算(S7)、多批共享一份期限 |
|
||||||
|
| 修改 | `tests/unit/test_ocr_client.py` | 取消结算恒 0(S6)、两个入口的期限与 per-call 校验 |
|
||||||
|
| 修改 | `tests/unit/test_client.py` | chat 期限命中、终态遥测行、已计费成功被丢弃、注入钟无关 |
|
||||||
|
| 修改 | `tests/unit/test_config.py` | 键未设/非法值 × env / 直接构造 / `dataclasses.replace` / 三个 `__init__` 直传 |
|
||||||
|
| 修改 | `tests/unit/test_openai_compat.py` | F1 四类取值 |
|
||||||
|
| 修改 | `tests/integration/test_redis_cross_connection.py` | 真实 Redis 上的取消结算契约(复用既有 `clients`/`_limiter`/`_client`/`ScriptedTransport`,**不改 Lua、不改契约套件**) |
|
||||||
|
| 修改 | `CHANGELOG.md`、`README.md`、`.env.example` | 新键、新异常、对外承诺三句话(§4 C4) |
|
||||||
|
| 新建 | `research-wiki/findings/2026-09-10-136-call-deadline-validation.md` | 红绿、命令、豁免索引,≤300 行 |
|
||||||
|
|
||||||
|
**不改**:`ports.py`(`Permit.settle` 签名与语义不动)、`middleware/admission.py`(`settle_and_release` 逐字不动)、`middleware/ratelimit.py`、`middleware/breaker.py`、`middleware/structured.py`、`middleware/cache.py`、`middleware/telemetry.py`、`telemetry/schema.py`(零新增列)、`backends/**`(含全部 Lua)、`sources.py`、`streaming.py`、`types.py`、`transports/monkey_ocr.py`(`:58` 不解析 `Retry-After`)、`tests/contracts/**`(复用现有 settle 契约,只跑不改)。若实施时发现必须突破本清单,先说明最小原因交父会话核定。
|
||||||
|
|
||||||
|
## 3. 跨任务接口(可执行定义,禁止占位)
|
||||||
|
|
||||||
|
### 3.1 `src/polygateway/deadline.py`
|
||||||
|
|
||||||
|
```python
|
||||||
|
def ensure_call_deadline(value: object, origin: str) -> float | None:
|
||||||
|
"""全装配路径共用的值域校验: None 或有限正数, 否则 ValueError(消息含 origin)。"""
|
||||||
|
if value is None:
|
||||||
|
return None
|
||||||
|
if isinstance(value, bool) or not isinstance(value, (int, float)):
|
||||||
|
raise ValueError(f"{origin} 必须是 None 或有限正数秒: {value!r}") # bool 先判
|
||||||
|
v = float(value)
|
||||||
|
if not math.isfinite(v) or v <= 0:
|
||||||
|
raise ValueError(f"{origin} 必须是有限正数秒: {value!r}") # NaN/inf/0/负
|
||||||
|
return v
|
||||||
|
|
||||||
|
|
||||||
|
async def with_call_deadline[T](aw: Awaitable[T], *, deadline_s: float | None, scope: str) -> T:
|
||||||
|
if deadline_s is None:
|
||||||
|
return await aw # 未启用: 不进上下文, 逐字旧路径
|
||||||
|
inner_timeout: BaseException | None = None
|
||||||
|
try:
|
||||||
|
async with asyncio.timeout(deadline_s) as cm:
|
||||||
|
try:
|
||||||
|
return await aw
|
||||||
|
except TimeoutError as exc:
|
||||||
|
inner_timeout = exc # 体内(含清理路径)自抛, 非本层期限
|
||||||
|
raise
|
||||||
|
except TimeoutError as exc:
|
||||||
|
if cm.expired() and exc is not inner_timeout:
|
||||||
|
raise CallDeadlineExceeded(scope=scope, deadline_s=deadline_s) from None
|
||||||
|
raise
|
||||||
|
```
|
||||||
|
|
||||||
|
三条实现红线:① **校验先于构造 awaitable**(否则非法值抛错时遗留未 await 协程 → `RuntimeWarning` + 资源不释放);② 只传**相对时长**,绝不把注入 `now` 加偏移换算成绝对截止时刻;③ 身份比较**不可退化**为只看 `cm.expired()`——探针 `/tmp/pgw_deadline_probe.py` E1/E2 实测:清理路径自抛的 `TimeoutError` 会被只看 `expired()` 的写法改标成 `CallDeadlineExceeded`;`__cause__` 启发式同样失效(内层 `asyncio.timeout` 的 `TimeoutError` 其 `__cause__` 也是 `CancelledError`)。本机制**不新增公共配置、不开后台任务、不改异常对象**。
|
||||||
|
|
||||||
|
### 3.2 `errors.py` 新类(唯一定义点)
|
||||||
|
|
||||||
|
```python
|
||||||
|
class CallDeadlineExceeded(PolyGatewayError):
|
||||||
|
"""调用方设定的整体期限到期; 不是网关不可用、也不是源故障。"""
|
||||||
|
|
||||||
|
def __init__(self, *, scope: str, deadline_s: float) -> None:
|
||||||
|
super().__init__(f"{scope} 调用期限 {deadline_s}s 到期")
|
||||||
|
self.scope = scope
|
||||||
|
self.deadline_s = deadline_s
|
||||||
|
```
|
||||||
|
|
||||||
|
无 `retry_after_s`(期限到期不含"何时可再试",给 `0.0` 会按既定语义指示下游立刻重打饱和渠道);不进 `SCOPE_REASONS`;不属四分类。`__init__.py` 导出后,三个边界既有的 `except PolyGatewayError` 自动接住并写终态行——**遥测零改动**。
|
||||||
|
|
||||||
|
### 3.3 配置(`config.py`)
|
||||||
|
|
||||||
|
| 项 | 精确定义 |
|
||||||
|
| --- | --- |
|
||||||
|
| 键名 | `{SCOPE}__CALL_DEADLINE_S`(两段式,`"LLM__CALL_DEADLINE_S".split("__")` 长度 **2** ≠ 4,故 `_load_sources`(函数定义 `:380`,判据行 `:384`)天然跳过,**不必**加进 `_RESERVED_SEGMENTS`) |
|
||||||
|
| loader | `_load_call_deadline(scope, env)`:`found = _first(env, f"{scope}__CALL_DEADLINE_S")`;`None → None`;否则 `ensure_call_deadline(_cast(found[1], "float", found[0]), found[0])`。**用 `_first` 不用 `_require`**(`_require` 会把未设当成配置缺失报错 = 破坏性变更);**origin 传实际命中的 env 键名 `found[0]`**,不传 `"GatewaySettings.call_deadline_s"`——否则 env 里写 `LLM__CALL_DEADLINE_S=0` 的人会拿到一条指向字段名的错误,在多 scope 部署里无法定位是哪个键(`_cast` 只能接住“不是数字”,`0`/负/`inf` 会穿过它) |
|
||||||
|
| 字段 | `GatewaySettings` **末尾**追加 `call_deadline_s: float \| None = None`(有默认值,不扰动既有位置构造;`EmbeddingSettings.gateway` / `OcrSettings.gateway` 自动继承) |
|
||||||
|
| 守卫 | `_validate_call_deadline()` 加入 `__post_init__`(`:186-194`)末位,实现体是 `object.__setattr__(self, "call_deadline_s", ensure_call_deadline(self.call_deadline_s, "GatewaySettings.call_deadline_s"))`——盖住直接构造与 `dataclasses.replace` 两条路;env 路已在 loader 里拿真键名报过错,此处重跑对合法值是幂等空操作 |
|
||||||
|
| 装配透传 | `GatewaySettings.from_env`(`:349`)返回字典加 `call_deadline_s=_load_call_deadline(scope_u, env)`;`GatewayClient.from_settings`(`client.py:449`)传 `settings.call_deadline_s`;`EmbeddingClient.from_settings`(`embedding.py:584`)与 `OcrClient.from_settings`(`ocr.py:631`)传 `gw.call_deadline_s`;三个 `from_env` 签名不变 |
|
||||||
|
| 不耦合 | **不校验** `call_deadline_s` 与 `timeout_s`/`stall_window_s` 的大小关系:期限短于单次超时是调用方的合法选择 |
|
||||||
|
|
||||||
|
### 3.4 三个接入点与 per-call 透传(唯一三处)
|
||||||
|
|
||||||
|
| 文件:行 | 改后 |
|
||||||
|
| --- | --- |
|
||||||
|
| `client.py:398` | `response = await with_call_deadline(self._handler(request), deadline_s=deadline, scope=self._scope)` |
|
||||||
|
| `embedding.py:209` | 同款包住 `self._embed_all(...)`(**整次调用一份**,N 批共享) |
|
||||||
|
| `ocr.py:274` | 同款包住 `self._run(...)` |
|
||||||
|
|
||||||
|
- 三处均在既有 `try` 之内、`_CallContext` 创建之后 → 到期照走 `except PolyGatewayError → emit_terminal_once`。
|
||||||
|
- `StructuredMW._run_ladder`(`structured.py:67-98`)与 `_embed_batch`(`embedding.py:295`)**严禁**新建 scope:同级重问/分批共享同一份期限,否则期限被轮数放大 N 倍。
|
||||||
|
- 三个 `__init__` 追加 keyword-only `call_deadline_s: float | None = None`,函数体首行即 `self._call_deadline_s = ensure_call_deadline(call_deadline_s, "<Class>(call_deadline_s=...)")`;`client.py` 需自存 `self._scope = scope`(现只传给 RetryMW,未自存)。
|
||||||
|
- 四个公开方法(`chat`/`embed`/`recognize_text`/`parse_layout`)追加 keyword-only `call_deadline_s: float | None = None`;`None` = 继承装配值,正数 = 本次覆盖,**不提供"本次关闭"**。取值一行:`deadline = self._call_deadline_s if call_deadline_s is None else ensure_call_deadline(call_deadline_s, "<method>(call_deadline_s=...)")`,位置在既有输入校验之列、`_CallContext` 创建**之前**。
|
||||||
|
- OCR 两个入口经 `recognize_text`/`parse_layout` → `_call` 形参透传;`_call` 内的校验须与 `image` 校验同列(`ocr.py:266-269`),即仍在 `_CallContext`(`:272`)之前。
|
||||||
|
- `embedding.py:205-217` 的 `texts == []` 早返回在 `try` 之前,天然落在期限之外——保持原样,测试显式记一笔。
|
||||||
|
|
||||||
|
### 3.5 取消结算矩阵落到代码(设计 §6.3 的唯一实现形态)
|
||||||
|
|
||||||
|
`middleware/retry.py::_attempt`(对照现读行号)。人类只批准了“**取消路径**的结算口径”这一条,故 `actual` 初值**不得**改成 `est`——那会把语义泛化到一切未分类异常(`RuntimeError`、`KeyError` 逃逸),属未批准范围。改用**局部阶段变量**:
|
||||||
|
|
||||||
|
```python
|
||||||
|
actual = 0 # 初值不动:未分类异常逃逸时仍逐字走 1.3.5 语义
|
||||||
|
settlement_known = False # 局部变量:该刻库是否已算出确定结算(不进任何签名)
|
||||||
|
```
|
||||||
|
|
||||||
|
| 现状行 | 现状 | 改后 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| `:281` | `actual = 0` | **保持 `0`**,紧随一行新增 `settlement_known = False`(局部阶段变量,带注释:只服务于取消分支的兜底取值,不进任何签名) |
|
||||||
|
| `:288` | `await self._transport.complete(...)` | 不动(S3 的“端口已开始”窗口就是它未返回的那段) |
|
||||||
|
| `:301` | `actual = source.effective_est_tokens()`(usage 不可得) | 值不变,其后置 `settlement_known = True` |
|
||||||
|
| `:303` | `actual = result.prompt_tokens + result.completion_tokens` | 值不变,其后置 `settlement_known = True`(S4;**真实 usage 恰为 0 也算已知**,取消不得覆写) |
|
||||||
|
| `:311` `RequestRejectedError` | 隐式 0 | 分支**首句**:`actual = 0; settlement_known = True`(逐字保住 1.3.5,且取消落在本分支 await 中途仍得 0) |
|
||||||
|
| `:315` `ResultInvalidError` | 隐式 0 | 同上(已计费的坏结果仍退全款属另一族缺口,本版不动) |
|
||||||
|
| `:320` `CancelledError` | 不动 `actual` | **唯一**新增赋值点,且在本分支首句:`if not settlement_known: actual = source.effective_est_tokens()`;`is_probe` / `_emit` / `raise` 三行原样 |
|
||||||
|
| `:325` 失败分支入口 | `dead = isinstance(exc, SourceDeadError)` | 完成该分支原有同步分类/选源反馈后,紧贴第一个 `await record_failure` 之前插入 `actual = 0 if dead else source.effective_est_tokens(); settlement_known = True`;不得前移到同步分类之前改变其异常结算 |
|
||||||
|
| `:335-336` | `if not dead: actual = est` | **删除**(已上移);非取消路径的最终值与 1.3.5 逐字相同,只是算得更早 |
|
||||||
|
| `:339-341` | `finally: pacer.leave(); settle_and_release(permit, actual)` | 一字不动 |
|
||||||
|
|
||||||
|
**“确定结算”的界桩就是上表的 await 位置**:失败分支的结算决定前移到两个记账 await 之前,故**取消发生在已知 `SourceDead` 之后时保留那个既有的 `0`**(源已判死就不该继续占额度)。**为何必须靠位置而不能只靠标志位**:`except asyncio.CancelledError` 与 `except (SourceDeadError, TransientError)` 是**同级**分支,落在后者块内 await 上的取消**不会**被前者接住,直接穿到 `finally`——那一刻 `actual` 是什么就结什么,标志位没有机会被读到。故失败分支必须在其**第一个 await 之前**就把 `actual` 定死。本版**不把 S5 泛化成“一切失败按 `est` 结算”**;`est` 只是“取消且结算未定”这一格的兜底值。
|
||||||
|
|
||||||
|
`embedding.py::_attempt` 同构:`:345` 保持 `actual = 0` 并新增 `settlement_known = False`;`:358`/`:360` 值不变、其后置 `True`;`:377`(`RequestRejected`/`ResultInvalid` 合并分支)首句 `actual = 0; settlement_known = True`;`:408` 失败分支在 `dead` 之后、`record_failure`(`:414`)之前插入 `actual = 0 if dead else est; settlement_known = True` 并删掉 `:414-415` 的 `if not dead:` 赋值;`:392` 取消分支首句加同款条件赋值;`:429-430 finally` 不动。
|
||||||
|
|
||||||
|
`ocr.py:449 settle_and_release(permit, 0)` **保持 0**:OCR 无 token 是事实而非"未知"(`ocr.py:9` 既有声明),不得改成 est,也不引入 `settlement_known`。
|
||||||
|
|
||||||
|
**防越界回归(必带)**:假 transport 抛 `RuntimeError`(不属四分类、无 except 接住)→ `tpm_used == 0` 且异常原样上抛;真实 usage 恰为 0 的成功 → `tpm_used == 0`。两条把“不得扩到未批准语义”钉成可回归的断言。
|
||||||
|
|
||||||
|
共享状态与探针一律不变:`pacer.leave()`、`permit.release()`、`breaker.release_probe()`(`retry.py:321-322`、`ocr.py:410-411`、`embedding.py:393-394`)、`mark_progress` 的调用点、次数与顺序全部逐字保留;**不新增任何公开参数**。
|
||||||
|
|
||||||
|
### 3.6 F1:`Retry-After` 非有限值(`transports/openai_compat.py:111-119`)
|
||||||
|
|
||||||
|
签名改为 `_parse_retry_after(raw: str | None, *, source_name: str) -> float | None`——`source_name` 是**必填 keyword-only 参数**(私有模块内函数,不属公共面,故不给默认值;漏传即 `TypeError`);**唯一调用处**是 `_translate_429`(`:140`),改传 `source_name=source.name`(该函数已持有 `source`,不需新参数)。
|
||||||
|
|
||||||
|
判据与告警(按已批设计 §6.1 原文精确定义,不得自行扩大):
|
||||||
|
|
||||||
|
| 输入形态 | 返回 | 日志 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| `inf` / `-inf` / `1e999`(`float()` 成功且 `math.isinf(seconds)`) | `None` | **一条 `logger.warning`**,只写源名与判据词(如 `retry_after_not_finite`),**不拼接、不截断、不打印原始头字符串** |
|
||||||
|
| `nan` | `None` | **无告警**:沿用既有 `seconds > 0` 恒假的值语义,本版**不为它新增分支、不改判据顺序** |
|
||||||
|
| HTTP-date / 空串 / 负数 / 不可解析 | `None` | 无告警(429 风暴下逐次告警会淹掉真信号) |
|
||||||
|
| 有限正数 | 该值 | 无 |
|
||||||
|
|
||||||
|
实现上只在 `float()` 成功后、`seconds > 0` 之前插一段 `if math.isinf(seconds): warning; return None`;`_translate_429` 的分类、`backoff_delay` 的 `max(delay, retry_after)` 取大逻辑一字不动(设计 §6.2:不夹 `backoff_max_s`)。
|
||||||
|
|
||||||
|
## 4. 任务与提交点(4 个原子提交)
|
||||||
|
|
||||||
|
### T0:设计批准状态与本计划(本任务,无代码)
|
||||||
|
|
||||||
|
产出:设计文档状态改批准 + §6.3 结算矩阵 + §5.2 五形态;本计划。不提交代码、不动测试。
|
||||||
|
|
||||||
|
### T1 → 提交 1 `fix: settle cancelled attempts against the source estimate`
|
||||||
|
|
||||||
|
1. **先红**:按 §5 批次 C 写 S3/S7 用例(`test_retry.py`、`test_embedding.py`),确认失败信息是 `tpm_used == 0 != 400`(不是构造错误);同批写 S5-dead 与 S8 两条**防越界**用例(实现前应已绿,作回归锁)。
|
||||||
|
2. 改 `middleware/retry.py::_attempt` 与 `embedding.py::_attempt`(§3.5:初值保 0 + 局部 `settlement_known`,失败分支结算决定上移到两个 await 之前),`ocr.py` 只补注释不改值;**不新增任何公开参数、不改未分类异常路径**。
|
||||||
|
3. **后绿**:新用例通过;`pytest tests/unit -q` 全绿(S1/S4/S5-dead/S6/S8 回归断言在批次 C 内一并落地)。
|
||||||
|
4. 真实 Redis:`tests/integration/test_redis_cross_connection.py` 新增取消结算用例(§5 批次 F),跑 `pytest tests/integration/test_redis_cross_connection.py -q`。
|
||||||
|
5. 暂存路径:`src/polygateway/middleware/retry.py`、`src/polygateway/embedding.py`、`src/polygateway/ocr.py`、三个测试文件。
|
||||||
|
|
||||||
|
### T2 → 提交 2 `feat: add an optional per-call wall-clock deadline`
|
||||||
|
|
||||||
|
1. 新建 `deadline.py`(§3.1)、`errors.py` 新类(§3.2)、`__init__.py` 导出、`pyproject.toml` layers 一行。
|
||||||
|
2. `config.py` 四处(字段/loader/守卫/`from_env`)、三个 client 的构造参数 + 公开方法参数 + 包裹点 + `from_settings` 透传(§3.3/§3.4)。
|
||||||
|
3. 先红后绿顺序:批次 A(`test_deadline.py` 值域)→ 批次 B(五形态)→ 批次 D(三链路命中与覆盖面)→ 批次 E(配置四条路)。
|
||||||
|
4. 回归门:`pytest tests/unit tests/contracts -q` 全绿且**未改一行既有断言**;`make lint`(含 import-linter 新层)通过。
|
||||||
|
5. 暂存路径:`src/polygateway/deadline.py`、`errors.py`、`__init__.py`、`config.py`、`client.py`、`embedding.py`、`ocr.py`、`pyproject.toml`、`tests/unit/test_deadline.py` 及四个改动测试文件。
|
||||||
|
|
||||||
|
### T3 → 提交 3 `fix: ignore non-finite Retry-After hints`
|
||||||
|
|
||||||
|
1. 先红:`tests/unit/test_openai_compat.py` 加 `inf`/`-inf`/`1e999`/`nan`/空/负/HTTP-date 七例,并加一例漏传 `source_name` 的 `TypeError`(批次 G)。
|
||||||
|
2. 改 `_parse_retry_after`:加必填私有 kw `source_name`、加 `math.isinf` 判据与一条 warning,同步唯一调用处 `_translate_429`(`:140`);跑该文件与 `tests/unit -q`。
|
||||||
|
3. 暂存:`src/polygateway/transports/openai_compat.py`、`tests/unit/test_openai_compat.py`。
|
||||||
|
|
||||||
|
### T4 → 提交 4 `docs: document the optional call deadline and cancellation settlement`
|
||||||
|
|
||||||
|
1. `CHANGELOG.md` 未发布段:三句强制措辞——**期限治理的是等待、返回时刻 = 期限 + 清理耗时(实测 5–7 倍)**;**到期不等于未产出、未计费**;**不配置即保持 1.3.5 语义**(纯 429 序列仍可能长等、有限大 `Retry-After` 仍照睡)。另记取消结算口径变化:**仅当取消发生在“端口已开始、结算尚未确定”时**按 `est` 保留预扣(方向为宁多扣不空退);已知结算(含真实 usage 恰为 0、已判 `SourceDead` 的 `0`)不被覆写,**未分类异常仍按 `0`**。另列"`except GatewayUnavailableError` 接不住新异常"。
|
||||||
|
2. `README.md` **四处同步(缺一不可,按行号定位)**:① 能力表(`:10-22` 区间)新增一行“调用期限”,措辞用 §5.3 三句;② “### 4. 业务侧异常处理”示例(`:185-195`)——该段 `except GatewayUnavailableError` **接不住** `CallDeadlineExceeded`,必须加一条 `except CallDeadlineExceeded` 分支并注明它无 `retry_after_s`;③ “哪些异常会到达调用方”表(`:466-476`)左列新增 `CallDeadlineExceeded` 行,并写明它**不属四分类、不属 `GatewayUnavailableError` 族**,只在显式配期限后才可能出现;④ 错误模型段补一句**有限大 `Retry-After` 残留**(能力表 `:15` 写的“尊重 Retry-After”仍成立:库不夹 `backoff_max_s`,服务端给 3600s 就睡 3600s,唯一制约手段是本版的调用期限;`inf`/`1e999` 自 1.3.6 起按无提示处理)。另:`.env.example` 在 `LLM__CIRCUIT_OPEN`(`:71`)之后加注释行 `# LLM__CALL_DEADLINE_S=`(缺省不启用,说明其治理对象是等待)。
|
||||||
|
3. `research-wiki/findings/2026-09-10-136-call-deadline-validation.md`:红绿证据、命令与退出码、豁免索引。
|
||||||
|
4. 独立验证(全新上下文 verifier)与整分支审查在本提交前完成;版本号与 wiki 同步留给发布清单(本计划不 bump、不发布)。
|
||||||
|
|
||||||
|
## 5. 测试矩阵 → 任务映射
|
||||||
|
|
||||||
|
**测试设施复用与"哪一份副本"的硬性核对**(历史坑,动手前必须核对):
|
||||||
|
|
||||||
|
| 事实 | 证据 | 纪律 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| `tests/unit/test_retry.py:35`、`tests/unit/test_embedding.py:176` 从 `tests.contracts.conftest` import `FakeClock` | 现读 | 改这一份即影响契约与两个单测文件 |
|
||||||
|
| `tests/unit/test_ocr_client.py:344` **自带一份同名 `FakeClock`** | 现读 | OCR 用例只吃这一份;给 OCR 加期限用例时不得误改 contracts 那份并以为生效 |
|
||||||
|
| `tests/unit/test_embedding.py:177` 从 `tests.unit.test_backpressure` import `BoundedSleep` | 现读 | 复用它做"轮询次数有界"断言,不新造 |
|
||||||
|
| 无 `tests/conftest.py` / `tests/unit/conftest.py` | `ls` 实测 | 新 fixture 只能进各文件本地,或复用 `tests/contracts/conftest.py`(已被 unit 直接 import) |
|
||||||
|
|
||||||
|
**取消白箱的确定性纪律**:既有取消用例用 `await asyncio.sleep(0.05)` 撞窗口(`test_retry.py:474`、`test_ocr_client.py:365`)——新用例**不得**沿用。做法:给 `test_retry.py::FakeTransport` 的 `"hang"` 分支加 `self.entered.set()`(构造期 `self.entered = asyncio.Event()`,3.10+ 不绑定 loop),用例 `await transport.entered.wait()` 后再 `task.cancel()`;embedding/OCR 的 `ScriptedEmbedTransport`/`ScriptedOcrTransport` 同款加一个 `entered`。既有用例不动。
|
||||||
|
|
||||||
|
| 批次 | 断言(→ 任务) | 落点 |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| A 值域 | `None` 通过;`0`/负/`nan`/`inf`/`"1"`/`True`(bool 不得当 1 秒)/`object()` → `ValueError` 且消息含 origin(→T2) | `tests/unit/test_deadline.py` |
|
||||||
|
| B 形态区分 | ①到期 → `CallDeadlineExceeded`(`scope`/`deadline_s` 正确);②未到期内层自抛 `TimeoutError` → 原样上抛;③**到期后清理自抛 `TimeoutError`** → 原样上抛且**断言不是** `CallDeadlineExceeded`(钉住身份比较);④外部 `task.cancel()`(先于/晚于到期各一例)→ `CancelledError`;⑤擦边成功 → 正常返回且 `task.cancelling() == 0`;⑥到期窗口内体内先抛领域异常(同步自旋构造)→ 上抛该异常,**不断言必为 deadline**;⑦`deadline_s=None` → 逐字旧路径(→T2) | `tests/unit/test_deadline.py`(真实 loop 时钟,期限 0.05s、体 0.3s,4–10× 余量,不标 slow) |
|
||||||
|
| C 取消结算 | S3:`tpm=1000, est_tokens=400`、transport `hang`、`entered` 后取消 → `tpm_used == 400`(**红→绿核心**)且 `inflight == 0`;S4:假 gate 在 `record_success` 处 `set()` 后挂起 → 取消 → `tpm_used == 15`(真实 usage 未被覆盖);S1:熔断开路使 `pick` 走 `settle_and_release(permit, 0)` → `tpm_used == 0`;S5-dead:假 gate 在 `SourceDead` 的 `record_failure` 处挂起 → 取消 → `tpm_used == 0`(**不得**变 `est`);S5-transient:同位置但瞬时失败 → `tpm_used == est`;S6:OCR 源 `tpm=600`、transport `hang` → 取消 → `tpm_used == 0`;S7:embedding 同 S3;**S8 防越界**:假 transport 抛 `RuntimeError` → `tpm_used == 0` 且异常原样上抛;真实 usage 恰为 0 的成功 → `tpm_used == 0`;四条既有路径(成功/`SourceDead`/`RequestRejected`/`ResultInvalid`)结算值逐字不变(→T1) | `test_retry.py`、`test_embedding.py`、`test_ocr_client.py` |
|
||||||
|
| D 覆盖面 | 期限分别落在 ①退避 `sleep`(注入真 `asyncio.sleep`)②准入排队(配额满轮询)③结构化重问 ④embedding 多批 → 均抛 `CallDeadlineExceeded`;embedding 断言 **N 批共享一份**期限(总时长不随批数放大);`texts == []` 早返回不受期限影响(→T2) | `test_client.py`、`test_embedding.py`、`test_ocr_client.py` |
|
||||||
|
| D2 到期代价 | ①终态遥测:到期恰好一条 `event_kind='terminal_failure'`、`error_type='CallDeadlineExceeded'`,被取消的 attempt 行仍 `cancelled`,两行 `logical_call_id` 一致,列数不变;②清理不可越过 + 量化:假 permit/假 emitter 各注入已知 sleep → 返回时刻 ≈ 期限 + 已知清理时长(断言 > 期限的若干倍,不断言上界);③已计费成功被丢弃:假缓存后端 `set` 慢于期限 → 抛 deadline 且断言 transport **已成功调用一次**(→T2) | `test_client.py` |
|
||||||
|
| E 配置四条路 | 键未设 → `None`;非法值 × {env、`GatewaySettings(...)` 直接构造、`dataclasses.replace`、三个 client `__init__` 直传} 各一例 → `ValueError`;per-call 非法值抛错且**无 "coroutine was never awaited" 警告**(`pytest.warns` 反向断言 / `-W error::RuntimeWarning`);`call_deadline_s < timeout_s` 合法不报错;注入钟跳变 10^6 秒**不**触发期限,而期限触发时 `total_latency_ms` 仍取自注入钟(→T2) | `test_config.py`、`test_client.py` |
|
||||||
|
| F 真实 Redis | 复用 `tests/integration/test_redis_cross_connection.py` 的 `clients`/`_limiter`/`_client`/`ScriptedTransport(hang=True)`:源 `tpm=1000, est_tokens=400`,`inflight` 出现后取消 → `source_stats.tpm_used == 400` 且 `inflight == 0`。**不改 Lua、不改 `tests/contracts/`**;另跑既有 `pytest tests/contracts/test_limiter_contract.py -q`(memory+redis 双参数)证明后端算术未被触碰(→T1) | `tests/integration/test_redis_cross_connection.py` |
|
||||||
|
| G F1 | `inf`/`-inf`/`1e999` → `None` + 各一条 warning(断言日志**含源名、不含**原始字符串);`nan`/空/负/HTTP-date → `None` 且**无** warning(`nan` 仍走既有 `seconds > 0` 值语义,不新增分支);漏传 `source_name` → `TypeError`(钉住必填 kw);有限正数仍参与 `max(delay, retry_after)`;`insufficient_quota` 仍归 `SourceDead`(→T3) | `test_openai_compat.py` |
|
||||||
|
| H 未启用回归 | `call_deadline_s=None` 时 `tests/unit`、`tests/contracts` 全绿且**未改一行既有断言**(→T2 门) | 全套件 |
|
||||||
|
|
||||||
|
命令(全部 `conda run -n PolyGateway`,禁止接管道以免退出码失真):`pytest tests/unit -q`、`pytest tests/contracts -q`、`pytest tests/integration/test_redis_cross_connection.py -q`、`make lint`。真实网关 e2e 与 `-m slow` 变体本计划**不跑**,由发布清单第 4 步统一负责。
|
||||||
|
|
||||||
|
## 6. 阻塞矩阵与交接
|
||||||
|
|
||||||
|
| 触发条件 | 处置 |
|
||||||
|
| --- | --- |
|
||||||
|
| 需要新增本计划外的公共键/端口方法/遥测列 | **停下上报**(设计 §9 边界之外即未批准) |
|
||||||
|
| 批次 B③(清理期自抛 `TimeoutError`)在实现里无法确定性构造 | 改用假端口在 `except CancelledError` 内直接 `raise TimeoutError`(探针 D3 已证可复现);仍不可得则记入 findings 的豁免索引,不得删断言 |
|
||||||
|
| 批次 D2② 的量化断言在 CI 机器上抖动 | 只断言下界(返回时刻 > 期限 × 2),不断言上界;不得改成 `sleep` 猜测 |
|
||||||
|
| 真实 Redis 不可用(`REDIS_URL` 未配置) | 用例自动 skip;findings 必须显式记"未取得真实 Redis 证据",不得以 memory 结果冒充 |
|
||||||
|
| 发现 `ResultInvalid`/`RequestRejected` 退全款想顺手修 | **不修**(设计 §6.3 末段:未批准的另一族记账语义),登记为新 issue 交父会话 |
|
||||||
|
| 失败分支结算决定上移后发现某条既有用例变红 | 先 `systematic-debugging` 定根因;如确为语义变化(非取消路径的最终值应与 1.3.5 逐字相同)则**停下上报**——说明本项只改算得更早、不改算出什么 |
|
||||||
|
| 想把保守口径扩到未分类异常(`RuntimeError` 等) | **不扩**(人类仅批准取消路径);S8 回归用例就是这道锁,需要就另立 issue |
|
||||||
|
| issue #24 相关想法 | 一律不实现、不写进代码与文档 |
|
||||||
|
|
||||||
|
交接物:4 个提交、1 份 findings、CHANGELOG 未发布段。版本号 bump、tag、构建、上传 registry 与 wiki 同步**不在本计划内**,按 CLAUDE.md §4.4.1 另行执行。
|
||||||
|
|
||||||
|
## 7. 自审
|
||||||
|
|
||||||
|
| 检查 | 结论 |
|
||||||
|
| --- | --- |
|
||||||
|
| 路径/行号/签名是否可执行无 TBD | 是——所有接入点均现读行号(`retry.py:281/288/301/303/311/315/320/325/335/339`、`embedding.py:345/349/358/360/377/392/408/414/429`、`ocr.py:449`、`client.py:398`、`embedding.py:209`、`ocr.py:274`、`config.py:186/349/380-384`、`openai_compat.py:111-119/140`、`__init__.py:61-62`) |
|
||||||
|
| 是否复用而非重造 | 是——`settle_and_release`、`emit_terminal_once`、`claim_terminal`、`asyncio.timeout` 范式、FakeClock/BoundedSleep/ScriptedTransport、限流契约套件与真实 Redis 用例全部复用;新增仅 1 文件 + 1 异常类 + 1 配置键 |
|
||||||
|
| 是否有先失败后通过的证据点 | 是——T1 的 S3/S7、T2 的 A/B/D、T3 的 G 均先红 |
|
||||||
|
| 取消与降级铁律 | 未新增 `except Exception`;`CancelledError` 无新捕获点;期限未启用时不进任何上下文 |
|
||||||
|
| 反 gold-plating | `ResultInvalid`/`RequestRejected` 退全款、#24、`shield`、遥测列、Lua 一律不碰 |
|
||||||
|
| 残余诚实标注 | S3 是保守选择而非"已计费"的证明;**未分类异常(S8)仍按 `0` 退全款,属已知残留、本版不动**;清理期自抛 `TimeoutError` 时无终态行;跨 Python 版本仅 3.12.13 有探针实证——四条均已写进设计 §12,findings 需复述 |
|
||||||
@@ -10,6 +10,7 @@ from polygateway.config import EmbeddingSettings, GatewaySettings, OcrSettings
|
|||||||
from polygateway.embedding import EmbeddingClient
|
from polygateway.embedding import EmbeddingClient
|
||||||
from polygateway.errors import (
|
from polygateway.errors import (
|
||||||
AllSourcesExhausted,
|
AllSourcesExhausted,
|
||||||
|
CallDeadlineExceeded,
|
||||||
CircuitOpenError,
|
CircuitOpenError,
|
||||||
GatewayUnavailableError,
|
GatewayUnavailableError,
|
||||||
GovernanceBackendError,
|
GovernanceBackendError,
|
||||||
@@ -51,7 +52,7 @@ from polygateway.types import (
|
|||||||
ThinkingObservation,
|
ThinkingObservation,
|
||||||
)
|
)
|
||||||
|
|
||||||
__version__ = "1.3.5"
|
__version__ = "1.3.6"
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
"DEFAULT_PROFILES",
|
"DEFAULT_PROFILES",
|
||||||
@@ -59,6 +60,7 @@ __all__ = [
|
|||||||
"Effort",
|
"Effort",
|
||||||
"AllSourcesExhausted",
|
"AllSourcesExhausted",
|
||||||
"CallStats",
|
"CallStats",
|
||||||
|
"CallDeadlineExceeded",
|
||||||
"CircuitOpenError",
|
"CircuitOpenError",
|
||||||
"EmbeddingClient",
|
"EmbeddingClient",
|
||||||
"EmbeddingResponse",
|
"EmbeddingResponse",
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ from polygateway.backends.memory.breaker import InMemoryGate
|
|||||||
from polygateway.backends.memory.cache import InMemoryCache
|
from polygateway.backends.memory.cache import InMemoryCache
|
||||||
from polygateway.backends.memory.limiter import InMemoryLimiter
|
from polygateway.backends.memory.limiter import InMemoryLimiter
|
||||||
from polygateway.config import GatewaySettings
|
from polygateway.config import GatewaySettings
|
||||||
|
from polygateway.deadline import ensure_call_deadline, with_call_deadline
|
||||||
from polygateway.errors import PolyGatewayError
|
from polygateway.errors import PolyGatewayError
|
||||||
from polygateway.middleware.base import compose
|
from polygateway.middleware.base import compose
|
||||||
from polygateway.middleware.cache import CacheMW
|
from polygateway.middleware.cache import CacheMW
|
||||||
@@ -235,10 +236,15 @@ class GatewayClient:
|
|||||||
structured_strategy: StructuredOutputStrategy | None = None,
|
structured_strategy: StructuredOutputStrategy | None = None,
|
||||||
structured_escalation: StructuredOutputStrategy | None = None,
|
structured_escalation: StructuredOutputStrategy | None = None,
|
||||||
structured_max_retries: int = 1,
|
structured_max_retries: int = 1,
|
||||||
|
call_deadline_s: float | None = None,
|
||||||
now: Any = time.monotonic,
|
now: Any = time.monotonic,
|
||||||
sleep: Any = asyncio.sleep,
|
sleep: Any = asyncio.sleep,
|
||||||
rng: Any = random.random,
|
rng: Any = random.random,
|
||||||
) -> None:
|
) -> None:
|
||||||
|
# 入口即校: 装配错误当场报,不等到第一次调用才炸
|
||||||
|
self._call_deadline_s = ensure_call_deadline(
|
||||||
|
call_deadline_s, "GatewayClient(call_deadline_s=...)"
|
||||||
|
)
|
||||||
emitter = (
|
emitter = (
|
||||||
TelemetryEmitter(telemetry, scope=scope, pricing=pricing, text_cap=text_cap)
|
TelemetryEmitter(telemetry, scope=scope, pricing=pricing, text_cap=text_cap)
|
||||||
if telemetry is not None
|
if telemetry is not None
|
||||||
@@ -292,6 +298,8 @@ class GatewayClient:
|
|||||||
self._structured_available = structured_strategy is not None
|
self._structured_available = structured_strategy is not None
|
||||||
self._terminal = terminal # 内部引用: 装配自省/测试用
|
self._terminal = terminal # 内部引用: 装配自省/测试用
|
||||||
self._handler = compose(middlewares, terminal)
|
self._handler = compose(middlewares, terminal)
|
||||||
|
# 期限到期需要报出 scope(现之前只传给 RetryMW,未自存)
|
||||||
|
self._scope = scope
|
||||||
# 逻辑调用统计需要同一只注入钟(1.3.5);现之前只传给中间件未自存
|
# 逻辑调用统计需要同一只注入钟(1.3.5);现之前只传给中间件未自存
|
||||||
self._now = now
|
self._now = now
|
||||||
# 终态行由公开边界统一写出(T3),故边界也需持有 emitter
|
# 终态行由公开边界统一写出(T3),故边界也需持有 emitter
|
||||||
@@ -336,6 +344,7 @@ class GatewayClient:
|
|||||||
reasoning_effort: Effort | str | None = None,
|
reasoning_effort: Effort | str | None = None,
|
||||||
tenant_id: str | None = None,
|
tenant_id: str | None = None,
|
||||||
meta: Mapping[str, Any] | None = None,
|
meta: Mapping[str, Any] | None = None,
|
||||||
|
call_deadline_s: float | None = None,
|
||||||
) -> LLMResponse:
|
) -> LLMResponse:
|
||||||
"""一次治理调用(签名冻结,ARCH §5.2;与三项目 LLMProvider 协议兼容)。
|
"""一次治理调用(签名冻结,ARCH §5.2;与三项目 LLMProvider 协议兼容)。
|
||||||
|
|
||||||
@@ -351,6 +360,11 @@ class GatewayClient:
|
|||||||
`tenant_id` 与 `meta` 是调用方自定义维度,只进遥测、**不进缓存 key**
|
`tenant_id` 与 `meta` 是调用方自定义维度,只进遥测、**不进缓存 key**
|
||||||
(租户隔离由 `cache_namespace` 负责,ARCH §7.5);前者享有真实列待遇
|
(租户隔离由 `cache_namespace` 负责,ARCH §7.5);前者享有真实列待遇
|
||||||
(可挂 RLS、可进复合索引),后者是任意 KV 容器(issue #11)。
|
(可挂 RLS、可进复合索引),后者是任意 KV 容器(issue #11)。
|
||||||
|
|
||||||
|
`call_deadline_s` 是本次调用的墙钟硬边界(issue #22): `None` = 继承装配值,
|
||||||
|
正数 = 本次覆盖,**不提供"本次关闭"**。它治理的是**等待**: 到期抛
|
||||||
|
`CallDeadlineExceeded`,但到期**不等于未产出、未计费**——在途请求可能已发出、
|
||||||
|
已被上游计费,且清理仍在 `finally` 里完成,故返回时刻 = 期限 + 清理耗时。
|
||||||
"""
|
"""
|
||||||
if structured is not None and not self._structured_available:
|
if structured is not None and not self._structured_available:
|
||||||
raise ImportError(
|
raise ImportError(
|
||||||
@@ -377,6 +391,13 @@ class GatewayClient:
|
|||||||
else coerce_effort(reasoning_effort, origin="chat(reasoning_effort=...)")
|
else coerce_effort(reasoning_effort, origin="chat(reasoning_effort=...)")
|
||||||
)
|
)
|
||||||
validate_thinking_raw(sampling, effort=effort, wire=None, origin="chat overlay")
|
validate_thinking_raw(sampling, effort=effort, wire=None, origin="chat overlay")
|
||||||
|
# 期限取值与校验必须在创建 awaitable **之前**: 否则非法值抛错时会遗留
|
||||||
|
# 未 await 的协程(RuntimeWarning + 资源不释放)
|
||||||
|
deadline = (
|
||||||
|
self._call_deadline_s
|
||||||
|
if call_deadline_s is None
|
||||||
|
else ensure_call_deadline(call_deadline_s, "chat(call_deadline_s=...)")
|
||||||
|
)
|
||||||
# 三项校验均已通过 → 进入统计边界(设计 §3: 输入校验异常在边界之外,保持原行为)
|
# 三项校验均已通过 → 进入统计边界(设计 §3: 输入校验异常在边界之外,保持原行为)
|
||||||
context = _CallContext(now=self._now)
|
context = _CallContext(now=self._now)
|
||||||
request = ChatRequest(
|
request = ChatRequest(
|
||||||
@@ -395,7 +416,9 @@ class GatewayClient:
|
|||||||
call_context=context,
|
call_context=context,
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
response = await self._handler(request)
|
response = await with_call_deadline(
|
||||||
|
self._handler(request), deadline_s=deadline, scope=self._scope
|
||||||
|
)
|
||||||
except PolyGatewayError as exc:
|
except PolyGatewayError as exc:
|
||||||
# 统计边界内的一切领域失败均尝试写一条终态行(1.3.5 设计 §6 I3),
|
# 统计边界内的一切领域失败均尝试写一条终态行(1.3.5 设计 §6 I3),
|
||||||
# 包括已有 attempt 错误行的 RequestRejected / ResultInvalid——两类行描述
|
# 包括已有 attempt 错误行的 RequestRejected / ResultInvalid——两类行描述
|
||||||
@@ -485,6 +508,7 @@ class GatewayClient:
|
|||||||
structured_strategy=strategy,
|
structured_strategy=strategy,
|
||||||
structured_escalation=escalation,
|
structured_escalation=escalation,
|
||||||
structured_max_retries=settings.structured_max_retries,
|
structured_max_retries=settings.structured_max_retries,
|
||||||
|
call_deadline_s=settings.call_deadline_s,
|
||||||
)
|
)
|
||||||
_mark_owned_components(client, limiter=limiter, breaker=breaker, telemetry=telemetry)
|
_mark_owned_components(client, limiter=limiter, breaker=breaker, telemetry=telemetry)
|
||||||
client._owns_cache = cache is None # 缓存后端可以是 None(backend=none),helper 会跳过
|
client._owns_cache = cache is None # 缓存后端可以是 None(backend=none),helper 会跳过
|
||||||
@@ -610,6 +634,7 @@ async def gather_bounded[T](aws: Iterable[Awaitable[T]], *, concurrency: int) ->
|
|||||||
"""有界并发 gather(D5 便利函数,替代 VT 手搓 semaphore+gather 样板)。
|
"""有界并发 gather(D5 便利函数,替代 VT 手搓 semaphore+gather 样板)。
|
||||||
|
|
||||||
语义与 `asyncio.gather` 默认一致: 结果保序、首个异常上抛;仅增加并发上限。
|
语义与 `asyncio.gather` 默认一致: 结果保序、首个异常上抛;仅增加并发上限。
|
||||||
|
期限计时从每次调用真正开始执行起算,信号量排队时长不在 `call_deadline_s` 之内。
|
||||||
"""
|
"""
|
||||||
if concurrency < 1:
|
if concurrency < 1:
|
||||||
raise ValueError("concurrency 必须 ≥ 1")
|
raise ValueError("concurrency 必须 ≥ 1")
|
||||||
|
|||||||
@@ -19,6 +19,7 @@ from typing import TYPE_CHECKING
|
|||||||
from dotenv import dotenv_values
|
from dotenv import dotenv_values
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
|
from polygateway.deadline import ensure_call_deadline
|
||||||
from polygateway.types import (
|
from polygateway.types import (
|
||||||
BackpressurePolicy,
|
BackpressurePolicy,
|
||||||
BreakerConfig,
|
BreakerConfig,
|
||||||
@@ -182,6 +183,11 @@ class GatewaySettings:
|
|||||||
pricing_path: str | None
|
pricing_path: str | None
|
||||||
structured_max_retries: int
|
structured_max_retries: int
|
||||||
lease_ttl_s: float
|
lease_ttl_s: float
|
||||||
|
# 一次逻辑调用的**可选**墙钟硬边界(issue #22)。缺省 None = 不启用,行为逐字
|
||||||
|
# 等于 1.3.5;有默认值故追加在末尾,不扰动既有位置构造。`EmbeddingSettings.gateway`
|
||||||
|
# 与 `OcrSettings.gateway` 自动继承。值域由 `_validate_call_deadline` 把关,
|
||||||
|
# 直接构造、`dataclasses.replace` 与 env 三条路一致
|
||||||
|
call_deadline_s: float | None = None
|
||||||
|
|
||||||
def __post_init__(self) -> None:
|
def __post_init__(self) -> None:
|
||||||
self._normalize()
|
self._normalize()
|
||||||
@@ -192,6 +198,7 @@ class GatewaySettings:
|
|||||||
self._validate_lease()
|
self._validate_lease()
|
||||||
self._validate_stall()
|
self._validate_stall()
|
||||||
self._validate_probe()
|
self._validate_probe()
|
||||||
|
self._validate_call_deadline()
|
||||||
|
|
||||||
def _normalize(self) -> None:
|
def _normalize(self) -> None:
|
||||||
"""把 `from_env` 一直在做的规范化补到构造路上,两条路必须产出同一个值。
|
"""把 `from_env` 一直在做的规范化补到构造路上,两条路必须产出同一个值。
|
||||||
@@ -345,6 +352,19 @@ class GatewaySettings:
|
|||||||
f"timeout_s + {_PROBE_GRACE_S}({floor});调大 probe_ttl_s 或调小源的 timeout_s"
|
f"timeout_s + {_PROBE_GRACE_S}({floor});调大 probe_ttl_s 或调小源的 timeout_s"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def _validate_call_deadline(self) -> None:
|
||||||
|
"""期限值域守卫: 盖住直接构造与 `dataclasses.replace` 两条路(issue #22)。
|
||||||
|
|
||||||
|
env 路已在 `_load_call_deadline` 里带真实键名报过错, 此处对合法值是幂等空操作。
|
||||||
|
**不校验**它与 `timeout_s`/`stall_window_s` 的大小关系: 期限短于单次超时
|
||||||
|
是调用方的合法选择(要的就是“不让这次调用拖过 N 秒”)。
|
||||||
|
"""
|
||||||
|
object.__setattr__(
|
||||||
|
self,
|
||||||
|
"call_deadline_s",
|
||||||
|
ensure_call_deadline(self.call_deadline_s, "GatewaySettings.call_deadline_s"),
|
||||||
|
)
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_env(
|
def from_env(
|
||||||
cls,
|
cls,
|
||||||
@@ -373,6 +393,7 @@ class GatewaySettings:
|
|||||||
selector=_load_choice(env, f"{scope_u}__SELECTOR", _SELECTORS, "health_aware"),
|
selector=_load_choice(env, f"{scope_u}__SELECTOR", _SELECTORS, "health_aware"),
|
||||||
quota_full=_load_choice(env, f"{scope_u}__QUOTA_FULL", _QUOTA_FULL, "wait"),
|
quota_full=_load_choice(env, f"{scope_u}__QUOTA_FULL", _QUOTA_FULL, "wait"),
|
||||||
circuit_open=_load_choice(env, f"{scope_u}__CIRCUIT_OPEN", _CIRCUIT_OPEN, "fail_fast"),
|
circuit_open=_load_choice(env, f"{scope_u}__CIRCUIT_OPEN", _CIRCUIT_OPEN, "fail_fast"),
|
||||||
|
call_deadline_s=_load_call_deadline(scope_u, env),
|
||||||
**_load_pgw(env),
|
**_load_pgw(env),
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -671,6 +692,30 @@ def _load_lease_ttl(env: Mapping[str, str]) -> float:
|
|||||||
return float(_cast(found[1], "float", found[0])) if found else _DEFAULT_LEASE_TTL_S
|
return float(_cast(found[1], "float", found[0])) if found else _DEFAULT_LEASE_TTL_S
|
||||||
|
|
||||||
|
|
||||||
|
def _load_call_deadline(scope: str, env: Mapping[str, str]) -> float | None:
|
||||||
|
"""读 `{SCOPE}__CALL_DEADLINE_S`(issue #22);键未设即 None = 不启用。
|
||||||
|
|
||||||
|
用 `_first` 而非 `_require`: 后者会把"未设"当成配置缺失报错,对存量下游
|
||||||
|
就是破坏性变更。键名两段式(`split("__")` 长度 2 ≠ 4),故 `_load_sources`
|
||||||
|
天然跳过它,不必进 `_RESERVED_SEGMENTS`。
|
||||||
|
|
||||||
|
origin 传**实际命中的 env 键名**而非字段名: `_cast` 只接得住"不是数字",
|
||||||
|
`0`/负数/`inf` 会穿过它落到 `ensure_call_deadline`——那时报一条指向字段名的错误,
|
||||||
|
在多 scope 部署里无法定位是哪个键写错了。
|
||||||
|
|
||||||
|
Args:
|
||||||
|
scope: 已大写的 scope 名。
|
||||||
|
env: 已合并的环境映射。
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
一次逻辑调用的墙钟期限(秒);键未设或为空串时返回 None(不启用)。
|
||||||
|
"""
|
||||||
|
found = _first(env, f"{scope}__CALL_DEADLINE_S")
|
||||||
|
if found is None:
|
||||||
|
return None
|
||||||
|
return ensure_call_deadline(_cast(found[1], "float", found[0]), found[0])
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True)
|
@dataclass(frozen=True)
|
||||||
class EmbeddingSettings:
|
class EmbeddingSettings:
|
||||||
"""Embedding scope 装配配置(M2 §7): 复用 GatewaySettings + embedding 专用键。
|
"""Embedding scope 装配配置(M2 §7): 复用 GatewaySettings + embedding 专用键。
|
||||||
|
|||||||
@@ -0,0 +1,80 @@
|
|||||||
|
"""一次逻辑调用的**可选**墙钟硬边界(issue #22;1.3.6 设计 §3 方案 A)。
|
||||||
|
|
||||||
|
只依赖标准库与 `errors.py`(依赖铁律最内层),供三个公开边界各包一次:
|
||||||
|
期限治理的是**等待**,不是"到期即无副作用"——在途请求可能已发出、已被上游
|
||||||
|
计费,清理照旧在 `finally` 完成,故返回时刻 = 期限 + 清理耗时。
|
||||||
|
|
||||||
|
缺省 `None` 时**完全不进上下文管理器**,行为逐字等于 1.3.5。
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import math
|
||||||
|
from typing import TYPE_CHECKING
|
||||||
|
|
||||||
|
from polygateway.errors import CallDeadlineExceeded
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
from collections.abc import Awaitable
|
||||||
|
|
||||||
|
|
||||||
|
def ensure_call_deadline(value: object, origin: str) -> float | None:
|
||||||
|
"""全装配路径共用的期限值域校验: `None` 或**有限正数秒**,否则当场 `ValueError`。
|
||||||
|
|
||||||
|
装配错误不属降级面(缺失/非法配置直接报错,不静默取默认值)。`bool` 必须先判:
|
||||||
|
`isinstance(True, int)` 为真,放行会让 `call_deadline_s=True` 变成"1 秒期限"
|
||||||
|
这种没人写得出来的意图。巨大 int(如 `10**400`)超出 float 值域,`float()` 会抛
|
||||||
|
`OverflowError`——它不是 `ValueError` 的子类,泄漏出去会绕过调用方的
|
||||||
|
`except ValueError`,故在此统一成同一种装配错误。
|
||||||
|
|
||||||
|
`origin` 写进消息,用于在多 scope 部署里定位到底是哪个键/哪个参数非法。
|
||||||
|
"""
|
||||||
|
if value is None:
|
||||||
|
return None
|
||||||
|
if isinstance(value, bool) or not isinstance(value, (int, float)):
|
||||||
|
raise ValueError(f"{origin} 必须是 None 或有限正数秒: {value!r}")
|
||||||
|
try:
|
||||||
|
seconds = float(value)
|
||||||
|
except OverflowError:
|
||||||
|
raise ValueError(f"{origin} 必须是有限正数秒: {value!r}") from None
|
||||||
|
if not math.isfinite(seconds) or seconds <= 0:
|
||||||
|
raise ValueError(f"{origin} 必须是有限正数秒: {value!r}")
|
||||||
|
return seconds
|
||||||
|
|
||||||
|
|
||||||
|
async def with_call_deadline[T](aw: Awaitable[T], *, deadline_s: float | None, scope: str) -> T:
|
||||||
|
"""给一个 awaitable 加一层可选期限;到期抛 `CallDeadlineExceeded`。
|
||||||
|
|
||||||
|
三条实现红线:
|
||||||
|
|
||||||
|
1. **校验先于构造 awaitable**——调用方必须先 `ensure_call_deadline`,否则非法值
|
||||||
|
抛错时会遗留未 await 的协程(`RuntimeWarning` + 资源不释放)。
|
||||||
|
2. 只用**相对时长**,绝不把注入的 `now` 换算成绝对截止时刻: 注入钟跳变
|
||||||
|
10^6 秒不该凭空触发期限。
|
||||||
|
3. 判据必须是**局部变量身份比较**,不可退化成只看 `cm.expired()`:
|
||||||
|
到期后清理路径自抛的 `TimeoutError` 也发生在 `expired()` 为真时,只看它
|
||||||
|
会把别人的超时改标成本层期限;`__cause__` 启发式同样失效(内层
|
||||||
|
`asyncio.timeout` 抛出的 `TimeoutError` 其 `__cause__` 也是 `CancelledError`)。
|
||||||
|
|
||||||
|
不新增后台任务、不 `shield`、不改异常对象:外部取消照常以 `CancelledError` 穿透。
|
||||||
|
"""
|
||||||
|
if deadline_s is None:
|
||||||
|
# 未启用: 不进上下文管理器,逐字走 1.3.5 旧路径
|
||||||
|
return await aw
|
||||||
|
# 体内(含清理路径)自抛的 TimeoutError 的**身份**,唯一可靠的区分依据
|
||||||
|
inner_timeout: BaseException | None = None
|
||||||
|
# 先建对象再进上下文: `as cm` 只在 `__aenter__` 返回后才绑定, 而 except 块无条件
|
||||||
|
# 读 `cm`——进入阶段一旦抛 TimeoutError 就会变成 NameError 掩盖真实错误
|
||||||
|
cm = asyncio.timeout(deadline_s)
|
||||||
|
try:
|
||||||
|
async with cm:
|
||||||
|
try:
|
||||||
|
return await aw
|
||||||
|
except TimeoutError as exc:
|
||||||
|
inner_timeout = exc
|
||||||
|
raise
|
||||||
|
except TimeoutError as exc:
|
||||||
|
if cm.expired() and exc is not inner_timeout:
|
||||||
|
raise CallDeadlineExceeded(scope=scope, deadline_s=deadline_s) from None
|
||||||
|
raise
|
||||||
@@ -28,6 +28,7 @@ from loguru import logger
|
|||||||
|
|
||||||
from polygateway.client import _aclose_component, _telemetry_status_of
|
from polygateway.client import _aclose_component, _telemetry_status_of
|
||||||
from polygateway.config import EmbeddingSettings
|
from polygateway.config import EmbeddingSettings
|
||||||
|
from polygateway.deadline import ensure_call_deadline, with_call_deadline
|
||||||
from polygateway.errors import (
|
from polygateway.errors import (
|
||||||
AllSourcesExhausted,
|
AllSourcesExhausted,
|
||||||
GovernanceBackendError,
|
GovernanceBackendError,
|
||||||
@@ -112,6 +113,7 @@ class EmbeddingClient:
|
|||||||
batch_size: int,
|
batch_size: int,
|
||||||
normalize: bool = False,
|
normalize: bool = False,
|
||||||
expected_dim: int | None = None,
|
expected_dim: int | None = None,
|
||||||
|
call_deadline_s: float | None = None,
|
||||||
now: Callable[[], float] = time.monotonic,
|
now: Callable[[], float] = time.monotonic,
|
||||||
sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
|
sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
|
||||||
rng: Callable[[], float] = random.random,
|
rng: Callable[[], float] = random.random,
|
||||||
@@ -120,6 +122,10 @@ class EmbeddingClient:
|
|||||||
raise ValueError("batch_size 必须 ≥ 1")
|
raise ValueError("batch_size 必须 ≥ 1")
|
||||||
if expected_dim is not None and expected_dim < 1:
|
if expected_dim is not None and expected_dim < 1:
|
||||||
raise ValueError("expected_dim 必须 ≥ 1")
|
raise ValueError("expected_dim 必须 ≥ 1")
|
||||||
|
# 入口即校: 装配错误当场报,不等到第一次调用才炸
|
||||||
|
self._call_deadline_s = ensure_call_deadline(
|
||||||
|
call_deadline_s, "EmbeddingClient(call_deadline_s=...)"
|
||||||
|
)
|
||||||
self._scope = scope
|
self._scope = scope
|
||||||
# embed payload 硬编码 {model, input},带 extra_body 的源必须先剥离,
|
# embed payload 硬编码 {model, input},带 extra_body 的源必须先剥离,
|
||||||
# 否则遥测会记录一个从未发出的采样参数(issue #4 决策 G)
|
# 否则遥测会记录一个从未发出的采样参数(issue #4 决策 G)
|
||||||
@@ -174,11 +180,16 @@ class EmbeddingClient:
|
|||||||
parent_call_id: str | None = None,
|
parent_call_id: str | None = None,
|
||||||
tenant_id: str | None = None,
|
tenant_id: str | None = None,
|
||||||
meta: Mapping[str, Any] | None = None,
|
meta: Mapping[str, Any] | None = None,
|
||||||
|
call_deadline_s: float | None = None,
|
||||||
) -> EmbeddingResponse:
|
) -> EmbeddingResponse:
|
||||||
"""一次治理 embedding 调用: 按 batch_size 切批,批间串行,全批合并返回。
|
"""一次治理 embedding 调用: 按 batch_size 切批,批间串行,全批合并返回。
|
||||||
|
|
||||||
`tenant_id` 与 `meta` 是调用方自定义维度,只进遥测(issue #11);它们属于
|
`tenant_id` 与 `meta` 是调用方自定义维度,只进遥测(issue #11);它们属于
|
||||||
本次调用而非某一批,故每批的遥测行都带同一份维度。
|
本次调用而非某一批,故每批的遥测行都带同一份维度。
|
||||||
|
|
||||||
|
`call_deadline_s` 是本次调用的墙钟硬边界(issue #22): `None` = 继承装配值。
|
||||||
|
**整次调用共享一份**——N 批串行跑在同一条期限内,不按批数放大 N 倍。
|
||||||
|
`texts == []` 的早返回在期限之外(零尝试,无等待可治)。
|
||||||
"""
|
"""
|
||||||
if not isinstance(texts, list) or any(not isinstance(t, str) for t in texts):
|
if not isinstance(texts, list) or any(not isinstance(t, str) for t in texts):
|
||||||
raise TypeError("texts 必须是 list[str](显式优于隐式,不收单条 str)")
|
raise TypeError("texts 必须是 list[str](显式优于隐式,不收单条 str)")
|
||||||
@@ -188,6 +199,12 @@ class EmbeddingClient:
|
|||||||
dimension_tenant_id, dimensions = validate_caller_dimensions(
|
dimension_tenant_id, dimensions = validate_caller_dimensions(
|
||||||
tenant_id, meta, origin="embed(tenant_id=..., meta=...)"
|
tenant_id, meta, origin="embed(tenant_id=..., meta=...)"
|
||||||
)
|
)
|
||||||
|
# 期限取值与校验必须在创建 awaitable **之前**(否则遗留未 await 的协程)
|
||||||
|
deadline = (
|
||||||
|
self._call_deadline_s
|
||||||
|
if call_deadline_s is None
|
||||||
|
else ensure_call_deadline(call_deadline_s, "embed(call_deadline_s=...)")
|
||||||
|
)
|
||||||
# 校验均已通过 → 进入统计边界(设计 §3.5: `texts` 类型与调用方维度校验之后)
|
# 校验均已通过 → 进入统计边界(设计 §3.5: `texts` 类型与调用方维度校验之后)
|
||||||
context = _CallContext(now=self._now)
|
context = _CallContext(now=self._now)
|
||||||
if not texts:
|
if not texts:
|
||||||
@@ -206,8 +223,12 @@ class EmbeddingClient:
|
|||||||
call_stats=context.snapshot(),
|
call_stats=context.snapshot(),
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
return await self._embed_all(
|
return await with_call_deadline(
|
||||||
texts, session_id, parent_call_id, dimension_tenant_id, dimensions, context
|
self._embed_all(
|
||||||
|
texts, session_id, parent_call_id, dimension_tenant_id, dimensions, context
|
||||||
|
),
|
||||||
|
deadline_s=deadline,
|
||||||
|
scope=self._scope,
|
||||||
)
|
)
|
||||||
except PolyGatewayError as exc:
|
except PolyGatewayError as exc:
|
||||||
await self._emit_terminal(
|
await self._emit_terminal(
|
||||||
@@ -343,6 +364,9 @@ class EmbeddingClient:
|
|||||||
call_id = str(uuid.uuid4())
|
call_id = str(uuid.uuid4())
|
||||||
started = self._now()
|
started = self._now()
|
||||||
actual = 0
|
actual = 0
|
||||||
|
# 局部阶段变量(与 RetryMW 同口径): 该刻库是否已算出确定结算。只服务于取消
|
||||||
|
# 分支的兜底取值,不进任何签名; 未分类异常逃逸时仍逐字走旧的全额退还。
|
||||||
|
settlement_known = False
|
||||||
# 登记在 transport 调用**之前**(同 RetryMW): 失败与取消的尝试也真的发出去了
|
# 登记在 transport 调用**之前**(同 RetryMW): 失败与取消的尝试也真的发出去了
|
||||||
context.register_attempt()
|
context.register_attempt()
|
||||||
try:
|
try:
|
||||||
@@ -358,6 +382,8 @@ class EmbeddingClient:
|
|||||||
actual = source.effective_est_tokens()
|
actual = source.effective_est_tokens()
|
||||||
else:
|
else:
|
||||||
actual = result.prompt_tokens
|
actual = result.prompt_tokens
|
||||||
|
# 真实 usage 恰为 0 也是已知事实, 后续取消不得改写成 est
|
||||||
|
settlement_known = True
|
||||||
await self._record_quietly(self._breaker.record_success(entry))
|
await self._record_quietly(self._breaker.record_success(entry))
|
||||||
await self._record_quietly(self._quota.mark_progress())
|
await self._record_quietly(self._quota.mark_progress())
|
||||||
latency_ms = int((self._now() - started) * 1000)
|
latency_ms = int((self._now() - started) * 1000)
|
||||||
@@ -375,6 +401,7 @@ class EmbeddingClient:
|
|||||||
)
|
)
|
||||||
return _BatchOutcome(result, source, call_id, latency_ms)
|
return _BatchOutcome(result, source, call_id, latency_ms)
|
||||||
except (RequestRejectedError, ResultInvalidError) as exc:
|
except (RequestRejectedError, ResultInvalidError) as exc:
|
||||||
|
actual, settlement_known = 0, True # 逐字保住 1.3.5 口径(本版不改这一族记账)
|
||||||
await self._gate_on_terminal(exc, entry)
|
await self._gate_on_terminal(exc, entry)
|
||||||
await self._emit(
|
await self._emit(
|
||||||
batch,
|
batch,
|
||||||
@@ -390,6 +417,9 @@ class EmbeddingClient:
|
|||||||
)
|
)
|
||||||
raise
|
raise
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
|
if not settlement_known:
|
||||||
|
# 端口已开始、结算未定: 保守保留预扣(设计 §6.3 S7)
|
||||||
|
actual = source.effective_est_tokens()
|
||||||
if entry.is_probe:
|
if entry.is_probe:
|
||||||
await self._record_quietly(self._breaker.release_probe(entry))
|
await self._record_quietly(self._breaker.release_probe(entry))
|
||||||
await self._emit(
|
await self._emit(
|
||||||
@@ -409,10 +439,11 @@ class EmbeddingClient:
|
|||||||
dead = isinstance(exc, SourceDeadError)
|
dead = isinstance(exc, SourceDeadError)
|
||||||
reason = _failure_reason(exc)
|
reason = _failure_reason(exc)
|
||||||
reasons[source.name] = reason
|
reasons[source.name] = reason
|
||||||
|
# 结算决定定死在本分支第一个 await 之前: 同级 except CancelledError 接不住
|
||||||
|
# 落在本块 await 上的取消,它直穿 finally。值与 1.3.5 逐字相同,只是算得更早。
|
||||||
|
actual = 0 if dead else source.effective_est_tokens()
|
||||||
|
settlement_known = True
|
||||||
await self._record_quietly(self._breaker.record_failure(entry, reason, dead))
|
await self._record_quietly(self._breaker.record_failure(entry, reason, dead))
|
||||||
if not dead:
|
|
||||||
# 保守: 失败请求可能已被网关计费(CHS 同款);与入场预扣同源取值
|
|
||||||
actual = source.effective_est_tokens()
|
|
||||||
await self._emit(
|
await self._emit(
|
||||||
batch,
|
batch,
|
||||||
source,
|
source,
|
||||||
@@ -624,6 +655,7 @@ class EmbeddingClient:
|
|||||||
batch_size=settings.batch_size,
|
batch_size=settings.batch_size,
|
||||||
normalize=settings.normalize,
|
normalize=settings.normalize,
|
||||||
expected_dim=settings.expected_dim,
|
expected_dim=settings.expected_dim,
|
||||||
|
call_deadline_s=gw.call_deadline_s,
|
||||||
)
|
)
|
||||||
_mark_owned_components(client, limiter=limiter, breaker=breaker, telemetry=telemetry)
|
_mark_owned_components(client, limiter=limiter, breaker=breaker, telemetry=telemetry)
|
||||||
return client
|
return client
|
||||||
|
|||||||
@@ -217,3 +217,23 @@ class GovernanceBackendError(GatewayUnavailableError):
|
|||||||
# 父类会把 message 覆写为 "{scope} 网关暂时不可用: {reason}",而各构造点
|
# 父类会把 message 覆写为 "{scope} 网关暂时不可用: {reason}",而各构造点
|
||||||
# 携带的诊断串(如"限流后端 try_acquire 失败: ...")是排障主线索,必须保住
|
# 携带的诊断串(如"限流后端 try_acquire 失败: ...")是排障主线索,必须保住
|
||||||
self.args = (message,)
|
self.args = (message,)
|
||||||
|
|
||||||
|
|
||||||
|
class CallDeadlineExceeded(PolyGatewayError): # noqa: N818 — 设计 §9 人类批准的公共名
|
||||||
|
"""调用方设定的整体调用期限到期; 不是网关不可用、也不是源故障。
|
||||||
|
|
||||||
|
刻意**不属**四分类、**不进** `SCOPE_REASONS`、**不继承** `GatewayUnavailableError`:
|
||||||
|
它描述的是调用方自己的耐心边界, 与"对方怎么了"正交——按四分类之一上报会让
|
||||||
|
下游的重试/换源/熔断逻辑对着一次本地超时做治理决策(库铁律「错误分类驱动」)。
|
||||||
|
|
||||||
|
也刻意**没有** `retry_after_s`: 期限到期不含"何时可再试"的信息, 给 `0.0`
|
||||||
|
会按既定语义指示下游立刻重打一条可能已经饱和的通道。
|
||||||
|
|
||||||
|
**到期不等于未产出、未计费**: 期限治理的是等待, 在途请求可能已经发出、
|
||||||
|
已被上游计费, 清理仍在 `finally` 里完成, 故返回时刻 = 期限 + 清理耗时。
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, *, scope: str, deadline_s: float) -> None:
|
||||||
|
super().__init__(f"{scope} 调用期限 {deadline_s}s 到期")
|
||||||
|
self.scope = scope
|
||||||
|
self.deadline_s = deadline_s
|
||||||
|
|||||||
@@ -279,6 +279,9 @@ class RetryMW:
|
|||||||
call_id = str(uuid.uuid4())
|
call_id = str(uuid.uuid4())
|
||||||
started = self._now()
|
started = self._now()
|
||||||
actual = 0
|
actual = 0
|
||||||
|
# 局部阶段变量: 该刻库是否已算出**确定**结算。只服务于取消分支的兜底取值,
|
||||||
|
# 不进任何签名/配置/遥测; 未分类异常逃逸时它无人读取, 故仍逐字走旧的全额退还。
|
||||||
|
settlement_known = False
|
||||||
# 登记在 transport 调用**之前**(1.3.5 设计 §4): 失败与取消的尝试同样
|
# 登记在 transport 调用**之前**(1.3.5 设计 §4): 失败与取消的尝试同样
|
||||||
# "真的打出去了",挪到成功之后会让诊断最需要看见的那几次从计数里消失。
|
# "真的打出去了",挪到成功之后会让诊断最需要看见的那几次从计数里消失。
|
||||||
# 上下文为 None = 库内现场构造的请求,跳过而不是报错
|
# 上下文为 None = 库内现场构造的请求,跳过而不是报错
|
||||||
@@ -301,6 +304,8 @@ class RetryMW:
|
|||||||
actual = source.effective_est_tokens()
|
actual = source.effective_est_tokens()
|
||||||
else:
|
else:
|
||||||
actual = result.prompt_tokens + result.completion_tokens
|
actual = result.prompt_tokens + result.completion_tokens
|
||||||
|
# 真实 usage 恰为 0 也是**已知事实**, 后续取消不得把它改写成 est
|
||||||
|
settlement_known = True
|
||||||
await self._record_quietly(self._breaker.record_success(entry))
|
await self._record_quietly(self._breaker.record_success(entry))
|
||||||
await self._record_quietly(self._quota.mark_progress())
|
await self._record_quietly(self._quota.mark_progress())
|
||||||
self._feed_outcome(source.name, ok=True)
|
self._feed_outcome(source.name, ok=True)
|
||||||
@@ -309,15 +314,20 @@ class RetryMW:
|
|||||||
await self._emit(request, source, call_id, started, response=response)
|
await self._emit(request, source, call_id, started, response=response)
|
||||||
return response
|
return response
|
||||||
except RequestRejectedError as exc:
|
except RequestRejectedError as exc:
|
||||||
|
actual, settlement_known = 0, True # 逐字保住 1.3.5: 坏请求全额退还
|
||||||
await self._on_rejected(exc, source, entry)
|
await self._on_rejected(exc, source, entry)
|
||||||
await self._emit(request, source, call_id, started, error=exc)
|
await self._emit(request, source, call_id, started, error=exc)
|
||||||
raise
|
raise
|
||||||
except ResultInvalidError as exc:
|
except ResultInvalidError as exc:
|
||||||
|
actual, settlement_known = 0, True # 同上, 本版不改这一族记账口径
|
||||||
# 坏结果 ≠ 坏服务: 熔断记成功但不计窗口样本,亦不喂健康分(M2.5 §3.1)
|
# 坏结果 ≠ 坏服务: 熔断记成功但不计窗口样本,亦不喂健康分(M2.5 §3.1)
|
||||||
await self._record_quietly(self._breaker.record_success(entry, count_attempt=False))
|
await self._record_quietly(self._breaker.record_success(entry, count_attempt=False))
|
||||||
await self._emit(request, source, call_id, started, error=exc)
|
await self._emit(request, source, call_id, started, error=exc)
|
||||||
raise
|
raise
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
|
if not settlement_known:
|
||||||
|
# 端口已开始、结算未定: 保守保留预扣(宁多扣不凭空退款, 见设计 §6.3 S3)
|
||||||
|
actual = source.effective_est_tokens()
|
||||||
if entry.is_probe:
|
if entry.is_probe:
|
||||||
await self._record_quietly(self._breaker.release_probe(entry))
|
await self._record_quietly(self._breaker.release_probe(entry))
|
||||||
await self._emit(request, source, call_id, started, error="cancelled")
|
await self._emit(request, source, call_id, started, error="cancelled")
|
||||||
@@ -330,10 +340,12 @@ class RetryMW:
|
|||||||
self._feed_outcome(source.name, ok=False)
|
self._feed_outcome(source.name, ok=False)
|
||||||
if reason == "rate_limited":
|
if reason == "rate_limited":
|
||||||
self._pacer.on_backpressure(source.name)
|
self._pacer.on_backpressure(source.name)
|
||||||
|
# 结算决定必须在本分支**第一个 await 之前**定死: 同级的 except CancelledError
|
||||||
|
# 接不住落在本块 await 上的取消,它会直穿 finally——那一刻 actual 是什么就结什么。
|
||||||
|
# 值与 1.3.5 逐字相同(dead 全额退、瞬时保留预扣),只是算得更早。
|
||||||
|
actual = 0 if dead else source.effective_est_tokens()
|
||||||
|
settlement_known = True
|
||||||
await self._record_quietly(self._breaker.record_failure(entry, reason, dead))
|
await self._record_quietly(self._breaker.record_failure(entry, reason, dead))
|
||||||
if not dead:
|
|
||||||
# 保守: 失败请求可能已被网关计费(CHS 同款);与入场预扣同源取值
|
|
||||||
actual = source.effective_est_tokens()
|
|
||||||
await self._emit(request, source, call_id, started, error=exc)
|
await self._emit(request, source, call_id, started, error=exc)
|
||||||
return _Failed(exc, immediate=dead)
|
return _Failed(exc, immediate=dead)
|
||||||
finally:
|
finally:
|
||||||
|
|||||||
+31
-2
@@ -23,6 +23,7 @@ from typing import TYPE_CHECKING, Any, Literal
|
|||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
from polygateway.client import _aclose_component, _telemetry_status_of
|
from polygateway.client import _aclose_component, _telemetry_status_of
|
||||||
|
from polygateway.deadline import ensure_call_deadline, with_call_deadline
|
||||||
from polygateway.errors import (
|
from polygateway.errors import (
|
||||||
AllSourcesExhausted,
|
AllSourcesExhausted,
|
||||||
GovernanceBackendError,
|
GovernanceBackendError,
|
||||||
@@ -115,10 +116,15 @@ class OcrClient:
|
|||||||
circuit_open: str = "fail_fast",
|
circuit_open: str = "fail_fast",
|
||||||
telemetry: TelemetryRecorder | None = None,
|
telemetry: TelemetryRecorder | None = None,
|
||||||
text_cap: int | None = None,
|
text_cap: int | None = None,
|
||||||
|
call_deadline_s: float | None = None,
|
||||||
now: Callable[[], float] = time.monotonic,
|
now: Callable[[], float] = time.monotonic,
|
||||||
sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
|
sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
|
||||||
rng: Callable[[], float] = random.random,
|
rng: Callable[[], float] = random.random,
|
||||||
) -> None:
|
) -> None:
|
||||||
|
# 入口即校: 装配错误当场报,不等到第一次调用才炸
|
||||||
|
self._call_deadline_s = ensure_call_deadline(
|
||||||
|
call_deadline_s, "OcrClient(call_deadline_s=...)"
|
||||||
|
)
|
||||||
self._scope = scope
|
self._scope = scope
|
||||||
# MonkeyOCR 只发 multipart 表单,带 extra_body 的源必须先剥离,否则
|
# MonkeyOCR 只发 multipart 表单,带 extra_body 的源必须先剥离,否则
|
||||||
# 遥测会记录一个从未发出的采样参数(issue #4 决策 G)
|
# 遥测会记录一个从未发出的采样参数(issue #4 决策 G)
|
||||||
@@ -171,10 +177,14 @@ class OcrClient:
|
|||||||
parent_call_id: str | None = None,
|
parent_call_id: str | None = None,
|
||||||
tenant_id: str | None = None,
|
tenant_id: str | None = None,
|
||||||
meta: Mapping[str, Any] | None = None,
|
meta: Mapping[str, Any] | None = None,
|
||||||
|
call_deadline_s: float | None = None,
|
||||||
) -> OcrTextResult:
|
) -> OcrTextResult:
|
||||||
"""一次治理文本转录(/ocr/text);text 空串 = 合法"无文字"。
|
"""一次治理文本转录(/ocr/text);text 空串 = 合法"无文字"。
|
||||||
|
|
||||||
`tenant_id` 与 `meta` 是调用方自定义维度,只进遥测(issue #11)。
|
`tenant_id` 与 `meta` 是调用方自定义维度,只进遥测(issue #11)。
|
||||||
|
|
||||||
|
`call_deadline_s` 是本次调用的墙钟硬边界(issue #22): `None` = 继承装配值。
|
||||||
|
它治理的是**等待**——到期不等于未产出,返回时刻 = 期限 + 清理耗时。
|
||||||
"""
|
"""
|
||||||
# 必须在进链路之前校验: 链路内的一切失败都被遥测层降级成 warning
|
# 必须在进链路之前校验: 链路内的一切失败都被遥测层降级成 warning
|
||||||
# (库铁律「遥测写失败降级不冒泡」),校验放下游等于没有校验
|
# (库铁律「遥测写失败降级不冒泡」),校验放下游等于没有校验
|
||||||
@@ -189,6 +199,7 @@ class OcrClient:
|
|||||||
parent_call_id,
|
parent_call_id,
|
||||||
dimension_tenant_id,
|
dimension_tenant_id,
|
||||||
dimensions,
|
dimensions,
|
||||||
|
call_deadline_s,
|
||||||
)
|
)
|
||||||
result = outcome.result
|
result = outcome.result
|
||||||
return OcrTextResult(
|
return OcrTextResult(
|
||||||
@@ -209,10 +220,13 @@ class OcrClient:
|
|||||||
parent_call_id: str | None = None,
|
parent_call_id: str | None = None,
|
||||||
tenant_id: str | None = None,
|
tenant_id: str | None = None,
|
||||||
meta: Mapping[str, Any] | None = None,
|
meta: Mapping[str, Any] | None = None,
|
||||||
|
call_deadline_s: float | None = None,
|
||||||
) -> OcrLayoutResult:
|
) -> OcrLayoutResult:
|
||||||
"""一次治理版面解析(/parse → ZIP);elements 空 = 合法"无元素"。
|
"""一次治理版面解析(/parse → ZIP);elements 空 = 合法"无元素"。
|
||||||
|
|
||||||
`tenant_id` 与 `meta` 是调用方自定义维度,只进遥测(issue #11)。
|
`tenant_id` 与 `meta` 是调用方自定义维度,只进遥测(issue #11)。
|
||||||
|
|
||||||
|
`call_deadline_s` 同 `recognize_text`(issue #22): `None` = 继承装配值。
|
||||||
"""
|
"""
|
||||||
# 校验早于链路,理由同 recognize_text;origin 标明方法名以便定位入口
|
# 校验早于链路,理由同 recognize_text;origin 标明方法名以便定位入口
|
||||||
dimension_tenant_id, dimensions = validate_caller_dimensions(
|
dimension_tenant_id, dimensions = validate_caller_dimensions(
|
||||||
@@ -226,6 +240,7 @@ class OcrClient:
|
|||||||
parent_call_id,
|
parent_call_id,
|
||||||
dimension_tenant_id,
|
dimension_tenant_id,
|
||||||
dimensions,
|
dimensions,
|
||||||
|
call_deadline_s,
|
||||||
)
|
)
|
||||||
result = outcome.result
|
result = outcome.result
|
||||||
return OcrLayoutResult(
|
return OcrLayoutResult(
|
||||||
@@ -262,17 +277,28 @@ class OcrClient:
|
|||||||
parent_call_id: str | None,
|
parent_call_id: str | None,
|
||||||
tenant_id: str | None,
|
tenant_id: str | None,
|
||||||
meta: dict[str, Any],
|
meta: dict[str, Any],
|
||||||
|
call_deadline_s: float | None = None,
|
||||||
) -> tuple[_AttemptOutcome, CallStats]:
|
) -> tuple[_AttemptOutcome, CallStats]:
|
||||||
if not isinstance(image, bytes):
|
if not isinstance(image, bytes):
|
||||||
raise TypeError("image 必须是 bytes(路径读取/批量拼帧留业务侧,D9)")
|
raise TypeError("image 必须是 bytes(路径读取/批量拼帧留业务侧,D9)")
|
||||||
if not image:
|
if not image:
|
||||||
raise ValueError("image 不能为空")
|
raise ValueError("image 不能为空")
|
||||||
|
# 期限校验与 `image` 校验同列(仍在 `_CallContext` 之前、创建 awaitable 之前)
|
||||||
|
deadline = (
|
||||||
|
self._call_deadline_s
|
||||||
|
if call_deadline_s is None
|
||||||
|
else ensure_call_deadline(call_deadline_s, f"{operation}(call_deadline_s=...)")
|
||||||
|
)
|
||||||
# M1 例外: `image` 校验在 `_call` 内而非公开方法,故上下文在该校验
|
# M1 例外: `image` 校验在 `_call` 内而非公开方法,故上下文在该校验
|
||||||
# **通过之后**创建——这样设计 §3 的"校验在统计边界外"对 OCR 才成立
|
# **通过之后**创建——这样设计 §3 的"校验在统计边界外"对 OCR 才成立
|
||||||
context = _CallContext(now=self._now)
|
context = _CallContext(now=self._now)
|
||||||
try:
|
try:
|
||||||
return await self._run(
|
return await with_call_deadline(
|
||||||
kind, operation, image, session_id, parent_call_id, tenant_id, meta, context
|
self._run(
|
||||||
|
kind, operation, image, session_id, parent_call_id, tenant_id, meta, context
|
||||||
|
),
|
||||||
|
deadline_s=deadline,
|
||||||
|
scope=self._scope,
|
||||||
)
|
)
|
||||||
except PolyGatewayError as exc:
|
except PolyGatewayError as exc:
|
||||||
await self._emit_terminal(
|
await self._emit_terminal(
|
||||||
@@ -446,6 +472,8 @@ class OcrClient:
|
|||||||
)
|
)
|
||||||
return _FailedAttempt(exc, immediate=dead)
|
return _FailedAttempt(exc, immediate=dead)
|
||||||
finally:
|
finally:
|
||||||
|
# OCR 的 0 token 是**事实**而非"未知"(设计 §6.3 S6): 故取消也恰恰结 0,
|
||||||
|
# 不引入 chat/embedding 那套 settlement_known 兜底。
|
||||||
await settle_and_release(permit, 0)
|
await settle_and_release(permit, 0)
|
||||||
|
|
||||||
async def _invoke(
|
async def _invoke(
|
||||||
@@ -668,6 +696,7 @@ class OcrClient:
|
|||||||
# OCR 行与 chat 行写同一张 llm_calls;漏传这一条,同表内就一半受控
|
# OCR 行与 chat 行写同一张 llm_calls;漏传这一条,同表内就一半受控
|
||||||
# 一半不受控(issue #12)
|
# 一半不受控(issue #12)
|
||||||
text_cap=gw.telemetry_text_cap,
|
text_cap=gw.telemetry_text_cap,
|
||||||
|
call_deadline_s=gw.call_deadline_s,
|
||||||
)
|
)
|
||||||
_mark_owned_components(client, limiter=limiter, breaker=breaker, telemetry=telemetry)
|
_mark_owned_components(client, limiter=limiter, breaker=breaker, telemetry=telemetry)
|
||||||
return client
|
return client
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ SSE 纯函数移植 VT `adapters/llm.py:51-124`;错误翻译移植 CHS
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
import json
|
||||||
|
import math
|
||||||
import re
|
import re
|
||||||
import time
|
import time
|
||||||
from dataclasses import replace
|
from dataclasses import replace
|
||||||
@@ -108,14 +109,30 @@ async def _iter_sse_deltas(
|
|||||||
# —— 错误翻译(CHS invokers.py 同款)——
|
# —— 错误翻译(CHS invokers.py 同款)——
|
||||||
|
|
||||||
|
|
||||||
def _parse_retry_after(raw: str | None) -> float | None:
|
def _parse_retry_after(raw: str | None, *, source_name: str) -> float | None:
|
||||||
"""解析 Retry-After 头;仅支持秒数形态,HTTP-date 返回 None(CHS 同款)。"""
|
"""解析 Retry-After 头;仅支持秒数形态,HTTP-date 返回 None(CHS 同款)。
|
||||||
|
|
||||||
|
**非有限值必须当作"无提示"**(issue F1): `"inf"` / `"1e999"` 能被 `float()`
|
||||||
|
成功解析,又能通过 `seconds > 0`,于是一路变成 `retry_after_s=inf`——而
|
||||||
|
`backoff_delay` 的 `max(delay, retry_after)` 取大之后就是一次**永不醒来**的
|
||||||
|
退避 sleep(库刻意不拿 `backoff_max_s` 去夹它,见设计 §6.2)。
|
||||||
|
|
||||||
|
`source_name` 是**必填** keyword-only 参数: 本函数是私有的,不给默认值,
|
||||||
|
漏传即 `TypeError`,免得将来新增调用点静默丢掉源标识(告警定位不到是哪个源)。
|
||||||
|
`nan` 不新增分支,沿用既有 `seconds > 0` 恒假的值语义。
|
||||||
|
"""
|
||||||
if raw is None:
|
if raw is None:
|
||||||
return None
|
return None
|
||||||
try:
|
try:
|
||||||
seconds = float(raw.strip())
|
seconds = float(raw.strip())
|
||||||
except ValueError:
|
except ValueError:
|
||||||
return None
|
return None
|
||||||
|
if math.isinf(seconds):
|
||||||
|
# 只写源名与判据词: 429 风暴下回显原始头会把日志淹掉,也无助于定位
|
||||||
|
logger.warning(
|
||||||
|
"{} 的 Retry-After 非有限值,按无提示处理(retry_after_not_finite)", source_name
|
||||||
|
)
|
||||||
|
return None
|
||||||
return seconds if seconds > 0 else None
|
return seconds if seconds > 0 else None
|
||||||
|
|
||||||
|
|
||||||
@@ -137,7 +154,7 @@ def _translate_429(
|
|||||||
)
|
)
|
||||||
return TransientError(
|
return TransientError(
|
||||||
compose_message(f"{source.name} 限速: 429", summary),
|
compose_message(f"{source.name} 限速: 429", summary),
|
||||||
retry_after_s=_parse_retry_after(headers.get("retry-after")),
|
retry_after_s=_parse_retry_after(headers.get("retry-after"), source_name=source.name),
|
||||||
**ctx,
|
**ctx,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -72,10 +72,13 @@ class ScriptedTransport:
|
|||||||
def __init__(self, hang: bool = False):
|
def __init__(self, hang: bool = False):
|
||||||
self.hang = hang
|
self.hang = hang
|
||||||
self.calls: list[str] = []
|
self.calls: list[str] = []
|
||||||
|
# 取消用例的确定性窗口(同 test_retry FakeTransport): 进入挂起即置位
|
||||||
|
self.entered = asyncio.Event()
|
||||||
|
|
||||||
async def complete(self, *, messages, source, stream, overlay, call_id, reasoning_effort):
|
async def complete(self, *, messages, source, stream, overlay, call_id, reasoning_effort):
|
||||||
self.calls.append(source.name)
|
self.calls.append(source.name)
|
||||||
if self.hang:
|
if self.hang:
|
||||||
|
self.entered.set()
|
||||||
await asyncio.Event().wait()
|
await asyncio.Event().wait()
|
||||||
return TransportResult(
|
return TransportResult(
|
||||||
content="ok",
|
content="ok",
|
||||||
@@ -225,7 +228,32 @@ async def test_cancel_in_flight_releases_lease(clients):
|
|||||||
assert (await limiter.source_stats("s1")).inflight == 0
|
assert (await limiter.source_stats("s1")).inflight == 0
|
||||||
|
|
||||||
|
|
||||||
# —— 掉线方向(fail-closed 集成证据)——
|
async def test_cancel_in_flight_keeps_the_reservation(clients):
|
||||||
|
"""1.3.6 §6.3 S3 在**真实 Redis** 上: 端口已开始、用量未知 → 保留 est 预扣。
|
||||||
|
|
||||||
|
内存后端与 Lua 后端的 `settle(delta)` 算术必须同口径——取消时凭空退款
|
||||||
|
在分布式部署下就是几个 worker 一起击穿 TPM 闸。本用例不改 Lua、不改契约套件。
|
||||||
|
"""
|
||||||
|
a_cli, _ = clients
|
||||||
|
scope = f"t{uuid4().hex[:8]}"
|
||||||
|
sources = [make_source(max_concurrency=1, tpm=1000, est_tokens=400)]
|
||||||
|
limiter = _limiter(a_cli, scope, sources, GlobalLimits(0, 0, 0))
|
||||||
|
transport = ScriptedTransport(hang=True)
|
||||||
|
client = _client(
|
||||||
|
scope,
|
||||||
|
sources,
|
||||||
|
limiter,
|
||||||
|
RedisGate(config=_CFG, redis=a_cli, scope=scope),
|
||||||
|
transport,
|
||||||
|
)
|
||||||
|
task = asyncio.create_task(client.chat([{"role": "user", "content": "hi"}]))
|
||||||
|
await transport.entered.wait()
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
stats = await limiter.source_stats("s1")
|
||||||
|
assert stats.tpm_used == 400 # 预扣保留, 不回退到 0
|
||||||
|
assert stats.inflight == 0
|
||||||
|
|
||||||
|
|
||||||
async def test_redis_down_admission_fails_closed():
|
async def test_redis_down_admission_fails_closed():
|
||||||
|
|||||||
@@ -1,8 +1,10 @@
|
|||||||
"""GatewayClient 装配与端到端(fake 后端 + MockTransport)测试(设计 §2.4)。"""
|
"""GatewayClient 装配与端到端(fake 后端 + MockTransport)测试(设计 §2.4)。"""
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
|
import gc
|
||||||
import json
|
import json
|
||||||
import sys
|
import sys
|
||||||
|
import warnings
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
@@ -10,10 +12,12 @@ import pytest
|
|||||||
|
|
||||||
from polygateway import (
|
from polygateway import (
|
||||||
AllSourcesExhausted,
|
AllSourcesExhausted,
|
||||||
|
CallDeadlineExceeded,
|
||||||
GatewayClient,
|
GatewayClient,
|
||||||
GatewaySettings,
|
GatewaySettings,
|
||||||
RequestRejectedError,
|
RequestRejectedError,
|
||||||
ResultInvalidError,
|
ResultInvalidError,
|
||||||
|
TransientError,
|
||||||
gather_bounded,
|
gather_bounded,
|
||||||
)
|
)
|
||||||
from polygateway.backends.memory.breaker import InMemoryGate
|
from polygateway.backends.memory.breaker import InMemoryGate
|
||||||
@@ -33,6 +37,10 @@ from polygateway.types import (
|
|||||||
SourceConfig,
|
SourceConfig,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# 复用 RetryMW 那份可编程 fake transport(含确定性 `entered` 窗口),不再造第二份;
|
||||||
|
# `tests/unit/test_backpressure.py:34` 已是同款复用
|
||||||
|
from tests.unit.test_retry import FakeTransport, _ok
|
||||||
|
|
||||||
_REPO = Path(__file__).resolve().parents[2]
|
_REPO = Path(__file__).resolve().parents[2]
|
||||||
|
|
||||||
_ENV = {
|
_ENV = {
|
||||||
@@ -1717,3 +1725,207 @@ class TestTerminalEmitDegradation:
|
|||||||
error=AllSourcesExhausted(scope="llm", reason="retry_exhausted", retry_after_s=1.0),
|
error=AllSourcesExhausted(scope="llm", reason="retry_exhausted", retry_after_s=1.0),
|
||||||
operation="chat",
|
operation="chat",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class _ClockJumpTransport:
|
||||||
|
"""假 transport: 只推进**注入钟**,真实墙钟几乎不走。
|
||||||
|
|
||||||
|
用于把"期限读哪只钟"与"统计读哪只钟"两件事分开断言。
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, clock, *, jump):
|
||||||
|
self._clock = clock
|
||||||
|
self._jump = jump
|
||||||
|
self.calls = []
|
||||||
|
|
||||||
|
async def complete(self, *, messages, source, stream, overlay, call_id, reasoning_effort):
|
||||||
|
self.calls.append(call_id)
|
||||||
|
self._clock.advance(self._jump)
|
||||||
|
return _ok()
|
||||||
|
|
||||||
|
|
||||||
|
class _SlowRecorder:
|
||||||
|
"""假 recorder: 每写一行真实等待一段,用于量化"清理不被期限截断"。"""
|
||||||
|
|
||||||
|
def __init__(self, delay=0.15):
|
||||||
|
self._delay = delay
|
||||||
|
self.rows = []
|
||||||
|
|
||||||
|
async def record_llm_call(self, **fields):
|
||||||
|
await asyncio.sleep(self._delay)
|
||||||
|
self.rows.append(fields)
|
||||||
|
|
||||||
|
|
||||||
|
class _SlowSetCache:
|
||||||
|
"""假缓存后端: `set` 慢于期限,用于构造"已产出、已计费的成功被丢弃"。"""
|
||||||
|
|
||||||
|
def __init__(self, delay=0.5):
|
||||||
|
self._delay = delay
|
||||||
|
self.data = {}
|
||||||
|
self.sets = 0
|
||||||
|
|
||||||
|
async def get(self, key):
|
||||||
|
return self.data.get(key)
|
||||||
|
|
||||||
|
async def set(self, key, value, ttl_s):
|
||||||
|
self.sets += 1
|
||||||
|
await asyncio.sleep(self._delay)
|
||||||
|
self.data[key] = value
|
||||||
|
|
||||||
|
|
||||||
|
class TestChatCallDeadline:
|
||||||
|
"""chat 链路的期限覆盖面与到期代价(计划 §5 批次 D/D2)。
|
||||||
|
|
||||||
|
真实事件循环时钟: 期限 0.05s 对被治理的等待(退避 5s、轮询 300s、慢 IO 0.5s)
|
||||||
|
有 10 倍以上余量,故不标 slow。
|
||||||
|
"""
|
||||||
|
|
||||||
|
_MSG = [{"role": "user", "content": "hi"}]
|
||||||
|
_DEADLINE = 0.05
|
||||||
|
|
||||||
|
def _rows(self, recorder, kind):
|
||||||
|
return [r for r in recorder.rows if r["event_kind"] == kind]
|
||||||
|
|
||||||
|
# —— 批次 D: 期限落点覆盖面 ——
|
||||||
|
|
||||||
|
async def test_deadline_fires_during_backoff_sleep(self):
|
||||||
|
"""退避 sleep 是等待的大头(429 序列可睡到小时级),期限必须能在它中间落地。"""
|
||||||
|
transport = FakeTransport([TransientError("boom", operation="chat"), _ok()])
|
||||||
|
async with _client(transport=transport, retry=RetryPolicy(3, 5.0, 30.0)) as client:
|
||||||
|
with pytest.raises(CallDeadlineExceeded) as exc:
|
||||||
|
await client.chat(self._MSG, call_deadline_s=self._DEADLINE)
|
||||||
|
assert exc.value.scope == "llm" and exc.value.deadline_s == self._DEADLINE
|
||||||
|
# 第二次尝试还压在 5s 退避里,期限确实落在 sleep 上而非 transport 上
|
||||||
|
assert len(transport.calls) == 1
|
||||||
|
|
||||||
|
async def test_deadline_fires_while_queued_for_quota(self):
|
||||||
|
"""准入排队(配额满轮询)是第二类长等待: 一次 transport 都没打出去也要能到期。"""
|
||||||
|
source = _source(tpm=1, est_tokens=1000) # 预扣量恒超本源 TPM → 六闸永不放行
|
||||||
|
limiter = InMemoryLimiter(
|
||||||
|
scope="llm", sources={source.name: source}, global_limits=GlobalLimits(0, 0, 0)
|
||||||
|
)
|
||||||
|
transport = FakeTransport([_ok()])
|
||||||
|
async with _client([source], transport=transport, limiter=limiter) as client:
|
||||||
|
with pytest.raises(CallDeadlineExceeded):
|
||||||
|
await client.chat(self._MSG, call_deadline_s=self._DEADLINE)
|
||||||
|
assert transport.calls == [] # 期限落在轮询里,尝试从未开始
|
||||||
|
|
||||||
|
async def test_deadline_fires_during_structured_re_ask(self):
|
||||||
|
"""结构化重问共享同一份期限: 阶梯不得按轮数各起一份,否则期限被放大 N 倍。
|
||||||
|
|
||||||
|
窗口用"第三轮挂起"构造而非 sleep 猜时长——前两轮瞬时返回坏 JSON,
|
||||||
|
期限只可能落在第三轮上,断言因此与机器负载无关。
|
||||||
|
"""
|
||||||
|
from pydantic import BaseModel
|
||||||
|
|
||||||
|
class Answer(BaseModel):
|
||||||
|
answer: int
|
||||||
|
|
||||||
|
transport = FakeTransport([_ok("not json at all"), _ok("not json at all"), "hang"])
|
||||||
|
async with _client(
|
||||||
|
transport=transport, structured_max_retries=5, structured_strategy=JsonRepairStrategy()
|
||||||
|
) as client:
|
||||||
|
with pytest.raises(CallDeadlineExceeded):
|
||||||
|
await client.chat(self._MSG, structured=Answer, call_deadline_s=0.05)
|
||||||
|
# 到期发生在第三轮: 期限确实跨过了两次重问,而不是在首轮就截断
|
||||||
|
assert len(transport.calls) == 3
|
||||||
|
|
||||||
|
# —— 批次 D2: 到期代价 ——
|
||||||
|
|
||||||
|
async def test_expiry_writes_one_terminal_row_and_a_cancelled_attempt(self):
|
||||||
|
"""到期恰好一条终态行 + 被取消的 attempt 行,两行同一 logical_call_id。"""
|
||||||
|
recorder = _MemoryRecorder()
|
||||||
|
transport = FakeTransport(["hang"])
|
||||||
|
async with _client(transport=transport, telemetry=recorder) as client:
|
||||||
|
with pytest.raises(CallDeadlineExceeded):
|
||||||
|
await client.chat(self._MSG, call_deadline_s=self._DEADLINE)
|
||||||
|
attempts = self._rows(recorder, "attempt")
|
||||||
|
terminals = self._rows(recorder, "terminal_failure")
|
||||||
|
assert len(terminals) == 1
|
||||||
|
assert terminals[0]["error_type"] == "CallDeadlineExceeded"
|
||||||
|
assert len(attempts) == 1 and attempts[0]["error"] == "cancelled"
|
||||||
|
assert attempts[0]["logical_call_id"] == terminals[0]["logical_call_id"]
|
||||||
|
|
||||||
|
# 零新增遥测列: 期限终态行的列集合与既有失败路径的终态行逐字相同
|
||||||
|
baseline = _MemoryRecorder()
|
||||||
|
async with _client(
|
||||||
|
handler=lambda request: httpx.Response(400, json={"error": {"message": "bad"}}),
|
||||||
|
telemetry=baseline,
|
||||||
|
) as client:
|
||||||
|
with pytest.raises(RequestRejectedError):
|
||||||
|
await client.chat(self._MSG)
|
||||||
|
assert set(terminals[0]) == set(self._rows(baseline, "terminal_failure")[0])
|
||||||
|
|
||||||
|
async def test_cleanup_is_not_cut_short_by_the_expiry(self):
|
||||||
|
"""返回时刻 = 期限 + 清理耗时: 只断下界(> 期限 × 2),不断上界。"""
|
||||||
|
recorder = _SlowRecorder(delay=0.15) # attempt 行与终态行各付一次
|
||||||
|
transport = FakeTransport(["hang"])
|
||||||
|
loop = asyncio.get_running_loop()
|
||||||
|
started = loop.time()
|
||||||
|
async with _client(transport=transport, telemetry=recorder) as client:
|
||||||
|
with pytest.raises(CallDeadlineExceeded):
|
||||||
|
await client.chat(self._MSG, call_deadline_s=self._DEADLINE)
|
||||||
|
elapsed = loop.time() - started
|
||||||
|
assert len(recorder.rows) == 2 # 清理照常写完两行,没被期限截断
|
||||||
|
assert elapsed > self._DEADLINE * 2, f"清理疑似被截断: {elapsed}s"
|
||||||
|
|
||||||
|
async def test_expiry_discards_a_success_that_was_already_billed(self):
|
||||||
|
"""到期 ≠ 未产出、未计费: transport 已成功一次,结果仍被丢弃。"""
|
||||||
|
transport = FakeTransport([_ok()])
|
||||||
|
cache = _SlowSetCache(delay=0.5) # 写缓存慢于期限 → 到期落在成功之后
|
||||||
|
async with _client(
|
||||||
|
transport=transport, cache=cache, cache_namespace="proj", cache_ttl_s=600
|
||||||
|
) as client:
|
||||||
|
with pytest.raises(CallDeadlineExceeded):
|
||||||
|
await client.chat(self._MSG, call_deadline_s=self._DEADLINE)
|
||||||
|
assert len(transport.calls) == 1 # 上游已经产出并计费
|
||||||
|
assert cache.sets == 1 and cache.data == {} # 结果既没回给调用方也没落缓存
|
||||||
|
|
||||||
|
# —— 批次 E: 注入钟与期限正交 ——
|
||||||
|
|
||||||
|
async def test_injected_clock_jump_does_not_trigger_the_deadline(self):
|
||||||
|
"""期限只认真实墙钟: 注入钟跳 10^6 秒也不该凭空到期(不换算绝对截止时刻)。"""
|
||||||
|
clock = _StatsClock()
|
||||||
|
transport = _ClockJumpTransport(clock, jump=1_000_000.0)
|
||||||
|
async with _client(transport=transport, now=clock) as client:
|
||||||
|
resp = await client.chat(self._MSG, call_deadline_s=5.0)
|
||||||
|
assert resp.content == "ok"
|
||||||
|
# 而统计仍逐字读注入钟(10^6 s = 10^9 ms),两只钟各司其职
|
||||||
|
assert resp.call_stats is not None
|
||||||
|
assert resp.call_stats.total_latency_ms == 1_000_000_000
|
||||||
|
|
||||||
|
async def test_expiry_latency_still_reads_the_injected_clock(self):
|
||||||
|
"""期限由真实钟触发,终态行的耗时仍取自注入钟(真实耗时只有几十毫秒)。"""
|
||||||
|
clock = _StatsClock()
|
||||||
|
recorder = _TickingRecorder(clock, tick=0.5)
|
||||||
|
transport = FakeTransport(["hang"])
|
||||||
|
async with _client(transport=transport, telemetry=recorder, now=clock) as client:
|
||||||
|
with pytest.raises(CallDeadlineExceeded):
|
||||||
|
await client.chat(self._MSG, call_deadline_s=self._DEADLINE)
|
||||||
|
terminal = [r for r in recorder.rows if r["event_kind"] == "terminal_failure"][0]
|
||||||
|
assert terminal["total_latency_ms"] == 500 # attempt 行那一次 tick,不是真实的 ~50ms
|
||||||
|
|
||||||
|
|
||||||
|
class TestChatCallDeadlineEntryGuards:
|
||||||
|
"""per-call 入口校验的两条硬红线(计划 §3.4/§5 批次 E)。"""
|
||||||
|
|
||||||
|
_MSG = [{"role": "user", "content": "hi"}]
|
||||||
|
|
||||||
|
async def test_illegal_per_call_value_leaves_no_un_awaited_coroutine(self):
|
||||||
|
"""校验先于构造 awaitable: 否则非法值抛错时遗留未 await 的协程(资源不释放)。"""
|
||||||
|
transport = FakeTransport([_ok()])
|
||||||
|
async with _client(transport=transport) as client:
|
||||||
|
with warnings.catch_warnings(record=True) as caught:
|
||||||
|
warnings.simplefilter("always")
|
||||||
|
with pytest.raises(ValueError, match=r"chat\(call_deadline_s"):
|
||||||
|
await client.chat(self._MSG, call_deadline_s=0)
|
||||||
|
gc.collect() # 未 await 的协程在回收时才发 RuntimeWarning
|
||||||
|
assert [w for w in caught if "never awaited" in str(w.message)] == []
|
||||||
|
assert transport.calls == []
|
||||||
|
|
||||||
|
async def test_per_call_none_inherits_the_assembled_deadline(self):
|
||||||
|
"""`None` = 继承装配值(不提供"本次关闭"): 装配了期限就照样到期。"""
|
||||||
|
transport = FakeTransport(["hang"])
|
||||||
|
async with _client(transport=transport, call_deadline_s=0.05) as client:
|
||||||
|
with pytest.raises(CallDeadlineExceeded):
|
||||||
|
await client.chat(self._MSG)
|
||||||
|
|||||||
@@ -1041,3 +1041,58 @@ def test_live_unknown_wire_assembly_is_local_only():
|
|||||||
GatewayClient.from_settings(
|
GatewayClient.from_settings(
|
||||||
dataclasses.replace(settings, sources=(source,)), registry=register_provider(mystery)
|
dataclasses.replace(settings, sources=(source,)), registry=register_provider(mystery)
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class TestCallDeadlineConfig:
|
||||||
|
"""`{SCOPE}__CALL_DEADLINE_S` 与三个 client 入口参数的值域四条路(issue #22)。"""
|
||||||
|
|
||||||
|
def test_key_unset_means_disabled(self):
|
||||||
|
assert GatewaySettings.from_env("LLM", env=_env()).call_deadline_s is None
|
||||||
|
|
||||||
|
def test_env_key_parsed(self):
|
||||||
|
s = GatewaySettings.from_env("LLM", env=_env(**{"LLM__CALL_DEADLINE_S": "30"}))
|
||||||
|
assert s.call_deadline_s == 30.0
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("bad", ["0", "-1", "nan", "inf", "abc"])
|
||||||
|
def test_env_illegal_value_reports_the_actual_key_name(self, bad):
|
||||||
|
"""origin 必须是实际命中的 env 键名,多 scope 部署里才定位得到是哪个键。"""
|
||||||
|
with pytest.raises(ValueError, match="LLM__CALL_DEADLINE_S"):
|
||||||
|
GatewaySettings.from_env("LLM", env=_env(**{"LLM__CALL_DEADLINE_S": bad}))
|
||||||
|
|
||||||
|
def test_direct_construction_is_guarded(self):
|
||||||
|
base = GatewaySettings.from_env("LLM", env=_env())
|
||||||
|
with pytest.raises(ValueError, match="GatewaySettings.call_deadline_s"):
|
||||||
|
dataclasses.replace(base, call_deadline_s=0)
|
||||||
|
|
||||||
|
def test_plain_constructor_call_is_guarded_too(self):
|
||||||
|
"""`dataclasses.replace` 与直接构造是两条路: 守卫在 `__post_init__` 才两条都盖住。
|
||||||
|
|
||||||
|
只在 `from_env` 里校验的话,直接 `GatewaySettings(...)` 装配的下游(测试/高级
|
||||||
|
注入路径,CLAUDE.md §4.5 的第二条装配路)会把非法期限一路带到第一次调用才炸。
|
||||||
|
"""
|
||||||
|
base = GatewaySettings.from_env("LLM", env=_env())
|
||||||
|
fields = {f.name: getattr(base, f.name) for f in dataclasses.fields(base)}
|
||||||
|
with pytest.raises(ValueError, match="GatewaySettings.call_deadline_s"):
|
||||||
|
GatewaySettings(**{**fields, "call_deadline_s": float("inf")})
|
||||||
|
# 合法值走同一条路不受影响(守卫对合法值是幂等空操作)
|
||||||
|
assert GatewaySettings(**{**fields, "call_deadline_s": 7}).call_deadline_s == 7.0
|
||||||
|
|
||||||
|
def test_replace_with_legal_value_is_idempotent(self):
|
||||||
|
base = GatewaySettings.from_env("LLM", env=_env())
|
||||||
|
assert dataclasses.replace(base, call_deadline_s=5).call_deadline_s == 5.0
|
||||||
|
|
||||||
|
def test_deadline_shorter_than_timeout_is_legal(self):
|
||||||
|
"""期限短于单次 timeout_s 是调用方的合法选择,不做跨字段耦合校验。"""
|
||||||
|
s = GatewaySettings.from_env("LLM", env=_env(**{"LLM__CALL_DEADLINE_S": "1"}))
|
||||||
|
assert s.call_deadline_s == 1.0 and s.sources[0].timeout_s == 120.0
|
||||||
|
|
||||||
|
def test_from_settings_propagates_to_client(self):
|
||||||
|
s = GatewaySettings.from_env("LLM", env=_env(**{"LLM__CALL_DEADLINE_S": "12"}))
|
||||||
|
assert GatewayClient.from_settings(s)._call_deadline_s == 12.0
|
||||||
|
|
||||||
|
def test_client_init_validates_at_entry(self):
|
||||||
|
"""三个 client 的 `__init__` 直传非法值也当场报错(不经 settings 那道守卫)。"""
|
||||||
|
from tests.unit.test_client import _client
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match=r"GatewayClient\(call_deadline_s"):
|
||||||
|
_client(call_deadline_s=0)
|
||||||
|
|||||||
@@ -0,0 +1,152 @@
|
|||||||
|
"""`deadline.py` 值域校验与五种形态区分测试(计划 §5 批次 A/B)。
|
||||||
|
|
||||||
|
用真实事件循环时钟(期限 0.05s、体 0.3s,4-10 倍余量),不标 slow:
|
||||||
|
被测对象是"哪一种 TimeoutError"的身份判据,注入钟无法覆盖 `asyncio.timeout`。
|
||||||
|
"""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import time
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from polygateway.deadline import ensure_call_deadline, with_call_deadline
|
||||||
|
from polygateway.errors import CallDeadlineExceeded
|
||||||
|
|
||||||
|
# —— 批次 A: 值域 ——
|
||||||
|
|
||||||
|
|
||||||
|
def test_ensure_call_deadline_accepts_none_and_positive():
|
||||||
|
assert ensure_call_deadline(None, "origin") is None
|
||||||
|
assert ensure_call_deadline(3, "origin") == 3.0
|
||||||
|
assert ensure_call_deadline(0.5, "origin") == 0.5
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"bad",
|
||||||
|
[0, 0.0, -1, -0.5, float("nan"), float("inf"), float("-inf"), "1", True, False, object(), []],
|
||||||
|
)
|
||||||
|
def test_ensure_call_deadline_rejects_out_of_range(bad):
|
||||||
|
with pytest.raises(ValueError) as exc:
|
||||||
|
ensure_call_deadline(bad, "GatewayClient(call_deadline_s=...)")
|
||||||
|
assert "GatewayClient(call_deadline_s=...)" in str(exc.value)
|
||||||
|
|
||||||
|
|
||||||
|
def test_ensure_call_deadline_rejects_huge_int_without_leaking_overflow():
|
||||||
|
"""超出 float 值域的巨大 int 也统一 ValueError,不泄漏 OverflowError。"""
|
||||||
|
with pytest.raises(ValueError) as exc:
|
||||||
|
ensure_call_deadline(10**400, "origin")
|
||||||
|
assert "origin" in str(exc.value)
|
||||||
|
|
||||||
|
|
||||||
|
# —— 批次 B: 五种形态 ——
|
||||||
|
|
||||||
|
|
||||||
|
async def test_deadline_expiry_raises_call_deadline_exceeded():
|
||||||
|
async def body():
|
||||||
|
await asyncio.sleep(0.3)
|
||||||
|
|
||||||
|
with pytest.raises(CallDeadlineExceeded) as exc:
|
||||||
|
await with_call_deadline(body(), deadline_s=0.05, scope="llm")
|
||||||
|
assert exc.value.scope == "llm"
|
||||||
|
assert exc.value.deadline_s == 0.05
|
||||||
|
|
||||||
|
|
||||||
|
async def test_inner_timeout_before_expiry_propagates_as_is():
|
||||||
|
async def body():
|
||||||
|
async with asyncio.timeout(0.01):
|
||||||
|
await asyncio.sleep(0.3)
|
||||||
|
|
||||||
|
with pytest.raises(TimeoutError) as exc:
|
||||||
|
await with_call_deadline(body(), deadline_s=5.0, scope="llm")
|
||||||
|
assert not isinstance(exc.value, CallDeadlineExceeded)
|
||||||
|
|
||||||
|
|
||||||
|
async def test_cleanup_timeout_after_expiry_is_not_relabelled():
|
||||||
|
"""到期后清理路径自抛 TimeoutError → 原样上抛(钉住身份比较,不看 expired())。"""
|
||||||
|
|
||||||
|
async def body():
|
||||||
|
try:
|
||||||
|
await asyncio.sleep(0.3)
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
raise TimeoutError("cleanup") from None
|
||||||
|
|
||||||
|
with pytest.raises(TimeoutError) as exc:
|
||||||
|
await with_call_deadline(body(), deadline_s=0.05, scope="llm")
|
||||||
|
assert not isinstance(exc.value, CallDeadlineExceeded)
|
||||||
|
assert str(exc.value) == "cleanup"
|
||||||
|
|
||||||
|
|
||||||
|
async def test_external_cancel_before_expiry_propagates_cancelled():
|
||||||
|
entered = asyncio.Event()
|
||||||
|
|
||||||
|
async def body():
|
||||||
|
entered.set()
|
||||||
|
await asyncio.sleep(0.3)
|
||||||
|
|
||||||
|
task = asyncio.create_task(with_call_deadline(body(), deadline_s=5.0, scope="llm"))
|
||||||
|
await entered.wait()
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
|
||||||
|
|
||||||
|
async def test_external_cancel_after_expiry_propagates_cancelled():
|
||||||
|
"""到期已在途、外部又取消 → 仍是 CancelledError(取消优先,不被改标)。"""
|
||||||
|
|
||||||
|
started = asyncio.Event()
|
||||||
|
|
||||||
|
async def body():
|
||||||
|
started.set()
|
||||||
|
try:
|
||||||
|
await asyncio.sleep(0.3)
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
await asyncio.sleep(0.2) # 清理期,期间遭外部取消
|
||||||
|
raise
|
||||||
|
|
||||||
|
task = asyncio.create_task(with_call_deadline(body(), deadline_s=0.05, scope="llm"))
|
||||||
|
await started.wait()
|
||||||
|
await asyncio.sleep(0.1) # 让期限先到期,进入清理
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
|
||||||
|
|
||||||
|
async def test_narrow_success_returns_value_without_pending_cancellation():
|
||||||
|
async def body():
|
||||||
|
await asyncio.sleep(0.01)
|
||||||
|
return "ok"
|
||||||
|
|
||||||
|
async def runner():
|
||||||
|
return await with_call_deadline(body(), deadline_s=0.2, scope="llm")
|
||||||
|
|
||||||
|
task = asyncio.create_task(runner())
|
||||||
|
assert await task == "ok"
|
||||||
|
assert task.cancelling() == 0
|
||||||
|
|
||||||
|
|
||||||
|
async def test_domain_error_inside_window_propagates():
|
||||||
|
"""计时器已触发但取消尚未投递的窗口内,体内先抛领域异常 → 原样上抛,期限静默让位。
|
||||||
|
|
||||||
|
忙等超过期限: 计时器回调已在 loop 上触发,但任务不挂起取消就投递不进来,
|
||||||
|
此刻体内同步抛出的领域异常必须原样逃逸(设计 §5.2 形态四)。
|
||||||
|
"""
|
||||||
|
|
||||||
|
class BoomError(RuntimeError):
|
||||||
|
pass
|
||||||
|
|
||||||
|
async def body():
|
||||||
|
end = time.monotonic() + 0.1
|
||||||
|
while time.monotonic() < end:
|
||||||
|
pass
|
||||||
|
raise BoomError("boom")
|
||||||
|
|
||||||
|
with pytest.raises(BoomError):
|
||||||
|
await with_call_deadline(body(), deadline_s=0.02, scope="llm")
|
||||||
|
|
||||||
|
|
||||||
|
async def test_none_deadline_takes_the_legacy_path():
|
||||||
|
async def body():
|
||||||
|
await asyncio.sleep(0.05)
|
||||||
|
return "ok"
|
||||||
|
|
||||||
|
assert await with_call_deadline(body(), deadline_s=None, scope="llm") == "ok"
|
||||||
@@ -13,6 +13,7 @@ import pytest
|
|||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
from polygateway.errors import (
|
from polygateway.errors import (
|
||||||
|
CallDeadlineExceeded,
|
||||||
RequestRejectedError,
|
RequestRejectedError,
|
||||||
ResultInvalidError,
|
ResultInvalidError,
|
||||||
SourceDeadError,
|
SourceDeadError,
|
||||||
@@ -197,6 +198,8 @@ class ScriptedEmbedTransport:
|
|||||||
def __init__(self, script):
|
def __init__(self, script):
|
||||||
self.script = list(script)
|
self.script = list(script)
|
||||||
self.calls = []
|
self.calls = []
|
||||||
|
# 取消用例的确定性窗口: 进入 hang 分支即置位, 不用 sleep 撞窗口
|
||||||
|
self.entered = asyncio.Event()
|
||||||
|
|
||||||
async def embed(self, *, texts, source, call_id):
|
async def embed(self, *, texts, source, call_id):
|
||||||
self.calls.append((source.name, list(texts), call_id))
|
self.calls.append((source.name, list(texts), call_id))
|
||||||
@@ -204,6 +207,7 @@ class ScriptedEmbedTransport:
|
|||||||
if isinstance(action, Exception):
|
if isinstance(action, Exception):
|
||||||
raise action
|
raise action
|
||||||
if action == "hang":
|
if action == "hang":
|
||||||
|
self.entered.set()
|
||||||
await asyncio.Event().wait()
|
await asyncio.Event().wait()
|
||||||
if action == "ok":
|
if action == "ok":
|
||||||
return _vec_for(texts)
|
return _vec_for(texts)
|
||||||
@@ -365,6 +369,20 @@ class TestEmbedGovernance:
|
|||||||
await task
|
await task
|
||||||
assert (await limiter.source_stats("e1")).inflight == 0
|
assert (await limiter.source_stats("e1")).inflight == 0
|
||||||
|
|
||||||
|
async def test_cancel_in_flight_keeps_the_reservation(self):
|
||||||
|
"""S7(与 chat 同口径): transport 在途被取消 → 用量未知 → 保留预扣而非退成 0。"""
|
||||||
|
transport = ScriptedEmbedTransport(["hang"])
|
||||||
|
client, limiter = _embed_client(
|
||||||
|
[_src(max_concurrency=1, tpm=1000, est_tokens=400)], [], transport=transport
|
||||||
|
)
|
||||||
|
task = asyncio.create_task(client.embed(["a"]))
|
||||||
|
await transport.entered.wait()
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
stats = await limiter.source_stats("e1")
|
||||||
|
assert stats.tpm_used == 400 and stats.inflight == 0
|
||||||
|
|
||||||
async def test_single_timeout_does_not_exhaust_stall_budget(self):
|
async def test_single_timeout_does_not_exhaust_stall_budget(self):
|
||||||
"""issue #8: 一次耗满 timeout 的尝试不得吃掉 stall 预算而使重试失效。
|
"""issue #8: 一次耗满 timeout 的尝试不得吃掉 stall 预算而使重试失效。
|
||||||
|
|
||||||
@@ -645,3 +663,70 @@ class TestEmbedLogicalCallStats:
|
|||||||
assert resp.call_stats.attempts == 0
|
assert resp.call_stats.attempts == 0
|
||||||
assert resp.call_stats.logical_call_id # 真实 ID,不是空串
|
assert resp.call_stats.logical_call_id # 真实 ID,不是空串
|
||||||
assert rec.rows == [] # 零遥测行
|
assert rec.rows == [] # 零遥测行
|
||||||
|
|
||||||
|
|
||||||
|
class _SlowEmbedTransport:
|
||||||
|
"""假 embedding transport: 每批真实耗时 `delay` 秒。
|
||||||
|
|
||||||
|
"N 批共享一份期限"只能用真实等待来证——`asyncio.timeout` 认的是事件循环
|
||||||
|
时钟,注入钟推不动它(计划 §5 批次 D)。
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, *, delay):
|
||||||
|
self._delay = delay
|
||||||
|
self.calls = []
|
||||||
|
|
||||||
|
async def embed(self, *, texts, source, call_id):
|
||||||
|
self.calls.append(list(texts))
|
||||||
|
await asyncio.sleep(self._delay)
|
||||||
|
return _vec_for(texts)
|
||||||
|
|
||||||
|
|
||||||
|
class TestEmbedCallDeadline:
|
||||||
|
"""embedding 的期限语义: 整次调用一份,空输入豁免(计划 §5 批次 D/E)。"""
|
||||||
|
|
||||||
|
async def test_one_deadline_is_shared_across_all_batches(self):
|
||||||
|
"""按批各起一份会让期限被批数放大 N 倍: 单批 0.05s 远小于期限 0.5s 时将永不到期。
|
||||||
|
|
||||||
|
余量刷到 10 倍(单批 0.05s vs 期限 0.5s): 要报假结论得单批慢 10 倍,
|
||||||
|
而不是机器抳一下就变色。
|
||||||
|
"""
|
||||||
|
transport = _SlowEmbedTransport(delay=0.05)
|
||||||
|
client, _ = _embed_client([_src()], [], batch_size=1, transport=transport)
|
||||||
|
loop = asyncio.get_running_loop()
|
||||||
|
started = loop.time()
|
||||||
|
with pytest.raises(CallDeadlineExceeded) as exc:
|
||||||
|
await client.embed([str(i) for i in range(20)], call_deadline_s=0.5)
|
||||||
|
elapsed = loop.time() - started
|
||||||
|
assert exc.value.scope == "embed"
|
||||||
|
# 按批计的话 20 批全都能跑完(根本不会抛),共享一份则跑不到头
|
||||||
|
assert 1 <= len(transport.calls) < 20
|
||||||
|
assert elapsed < 20 * 0.05, f"总时长疑似随批数放大: {elapsed}s"
|
||||||
|
|
||||||
|
async def test_empty_input_is_exempt_from_the_deadline(self):
|
||||||
|
"""`texts == []` 早返回在 try 之外(零尝试、无等待可治),再小的期限也不该拦它。"""
|
||||||
|
transport = ScriptedEmbedTransport([])
|
||||||
|
client, _ = _embed_client([_src()], [], transport=transport)
|
||||||
|
resp = await client.embed([], call_deadline_s=1e-6)
|
||||||
|
assert resp.vectors == []
|
||||||
|
assert resp.call_stats is not None and resp.call_stats.attempts == 0
|
||||||
|
assert transport.calls == []
|
||||||
|
|
||||||
|
async def test_the_same_tiny_deadline_does_fire_on_a_non_empty_input(self):
|
||||||
|
"""对照组: 上一条用的 1e-6 秒确实是会到期的值,豁免不是因为期限没生效。"""
|
||||||
|
transport = _SlowEmbedTransport(delay=0.05)
|
||||||
|
client, _ = _embed_client([_src()], [], transport=transport)
|
||||||
|
with pytest.raises(CallDeadlineExceeded):
|
||||||
|
await client.embed(["a"], call_deadline_s=1e-6)
|
||||||
|
|
||||||
|
async def test_illegal_per_call_value_is_rejected_at_the_entry(self):
|
||||||
|
"""per-call 非法值当场 ValueError,且消息指向 `embed(...)` 而非某个 env 键。"""
|
||||||
|
transport = ScriptedEmbedTransport([])
|
||||||
|
client, _ = _embed_client([_src()], [], transport=transport)
|
||||||
|
with pytest.raises(ValueError, match=r"embed\(call_deadline_s"):
|
||||||
|
await client.embed(["a"], call_deadline_s=0)
|
||||||
|
assert transport.calls == []
|
||||||
|
|
||||||
|
def test_illegal_constructor_value_is_rejected_at_assembly(self):
|
||||||
|
with pytest.raises(ValueError, match=r"EmbeddingClient\(call_deadline_s"):
|
||||||
|
_embed_client([_src()], [], call_deadline_s=-1)
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ from polygateway.backends.memory.breaker import InMemoryGate
|
|||||||
from polygateway.backends.memory.limiter import InMemoryLimiter
|
from polygateway.backends.memory.limiter import InMemoryLimiter
|
||||||
from polygateway.errors import (
|
from polygateway.errors import (
|
||||||
AllSourcesExhausted,
|
AllSourcesExhausted,
|
||||||
|
CallDeadlineExceeded,
|
||||||
CircuitOpenError,
|
CircuitOpenError,
|
||||||
RequestRejectedError,
|
RequestRejectedError,
|
||||||
ResultInvalidError,
|
ResultInvalidError,
|
||||||
@@ -60,6 +61,8 @@ class ScriptedOcrTransport:
|
|||||||
def __init__(self, script):
|
def __init__(self, script):
|
||||||
self.script = list(script)
|
self.script = list(script)
|
||||||
self.calls = []
|
self.calls = []
|
||||||
|
# 取消用例的确定性窗口: 进入 hang 分支即置位, 不用 sleep 撞窗口
|
||||||
|
self.entered = asyncio.Event()
|
||||||
|
|
||||||
async def _next(self, method, source, call_id):
|
async def _next(self, method, source, call_id):
|
||||||
self.calls.append((method, source.name, call_id))
|
self.calls.append((method, source.name, call_id))
|
||||||
@@ -67,6 +70,7 @@ class ScriptedOcrTransport:
|
|||||||
if isinstance(action, Exception):
|
if isinstance(action, Exception):
|
||||||
raise action
|
raise action
|
||||||
if action == "hang":
|
if action == "hang":
|
||||||
|
self.entered.set()
|
||||||
await asyncio.Event().wait()
|
await asyncio.Event().wait()
|
||||||
return _TEXT_OK if action == "text" else _LAYOUT_OK
|
return _TEXT_OK if action == "text" else _LAYOUT_OK
|
||||||
|
|
||||||
@@ -395,6 +399,20 @@ class TestCancellation:
|
|||||||
stats = await limiter.source_stats("m1")
|
stats = await limiter.source_stats("m1")
|
||||||
assert stats.inflight == 0 # permit 在 finally 释放
|
assert stats.inflight == 0 # permit 在 finally 释放
|
||||||
|
|
||||||
|
async def test_cancel_in_flight_still_settles_zero(self):
|
||||||
|
"""S6: OCR 的 0 token 是**事实**而非"未知", 取消也不得改成按 est 结算。"""
|
||||||
|
transport = ScriptedOcrTransport(["hang"])
|
||||||
|
client, limiter, _ = _client(
|
||||||
|
[_src(max_concurrency=1, tpm=1000, est_tokens=400)], [], transport=transport
|
||||||
|
)
|
||||||
|
task = asyncio.create_task(client.recognize_text(b"jpg"))
|
||||||
|
await transport.entered.wait()
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
stats = await limiter.source_stats("m1")
|
||||||
|
assert stats.tpm_used == 0 and stats.inflight == 0
|
||||||
|
|
||||||
|
|
||||||
class TestCheckHealth:
|
class TestCheckHealth:
|
||||||
class _HealthTransport(ScriptedOcrTransport):
|
class _HealthTransport(ScriptedOcrTransport):
|
||||||
@@ -652,3 +670,29 @@ class TestOcrLogicalCallStats:
|
|||||||
await client.recognize_text("not-bytes")
|
await client.recognize_text("not-bytes")
|
||||||
with pytest.raises(ValueError):
|
with pytest.raises(ValueError):
|
||||||
await client.recognize_text(b"")
|
await client.recognize_text(b"")
|
||||||
|
|
||||||
|
|
||||||
|
class TestOcrCallDeadline:
|
||||||
|
"""OCR 两个公开入口的期限与 per-call 校验(计划 §3.4/§5 批次 D/E)。"""
|
||||||
|
|
||||||
|
async def test_expiry_on_a_hanging_transport(self):
|
||||||
|
client, limiter, _ = _client([_src()], ["hang"])
|
||||||
|
with pytest.raises(CallDeadlineExceeded) as exc:
|
||||||
|
await client.recognize_text(b"jpg", call_deadline_s=0.05)
|
||||||
|
assert exc.value.scope == "ocr" and exc.value.deadline_s == 0.05
|
||||||
|
# 清理照常在 finally 完成: 在途计数必须归零(OCR 无 token,结算恒 0)
|
||||||
|
stats = await limiter.source_stats("m1")
|
||||||
|
assert stats.inflight == 0 and stats.tpm_used == 0
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("method", ["recognize_text", "parse_layout"])
|
||||||
|
async def test_illegal_per_call_value_names_the_entry_it_came_from(self, method):
|
||||||
|
"""两个入口各自报自己的名字: 多入口部署里才定位得到是哪次调用传错了。"""
|
||||||
|
transport = ScriptedOcrTransport([])
|
||||||
|
client, _, _ = _client([_src()], [], transport=transport)
|
||||||
|
with pytest.raises(ValueError, match=rf"{method}\(call_deadline_s"):
|
||||||
|
await getattr(client, method)(b"jpg", call_deadline_s=float("inf"))
|
||||||
|
assert transport.calls == []
|
||||||
|
|
||||||
|
def test_illegal_constructor_value_is_rejected_at_assembly(self):
|
||||||
|
with pytest.raises(ValueError, match=r"OcrClient\(call_deadline_s"):
|
||||||
|
_client([_src()], [], call_deadline_s=0)
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ from polygateway.errors import (
|
|||||||
SourceDeadError,
|
SourceDeadError,
|
||||||
TransientError,
|
TransientError,
|
||||||
)
|
)
|
||||||
|
from polygateway.middleware.retry import backoff_delay
|
||||||
from polygateway.middleware.telemetry import TelemetryEmitter
|
from polygateway.middleware.telemetry import TelemetryEmitter
|
||||||
from polygateway.pricing import ModelPrice, PricingTable
|
from polygateway.pricing import ModelPrice, PricingTable
|
||||||
from polygateway.providers import ProviderProfile, ThinkingWire, register_provider
|
from polygateway.providers import ProviderProfile, ThinkingWire, register_provider
|
||||||
@@ -21,9 +22,18 @@ from polygateway.transports._http_errors import summarize_body
|
|||||||
from polygateway.transports.openai_compat import (
|
from polygateway.transports.openai_compat import (
|
||||||
OpenAICompatTransport,
|
OpenAICompatTransport,
|
||||||
_iter_sse_deltas,
|
_iter_sse_deltas,
|
||||||
|
_parse_retry_after,
|
||||||
_sse_data_payload,
|
_sse_data_payload,
|
||||||
|
_translate_429,
|
||||||
|
)
|
||||||
|
from polygateway.types import (
|
||||||
|
ChatRequest,
|
||||||
|
Effort,
|
||||||
|
LLMResponse,
|
||||||
|
RetryPolicy,
|
||||||
|
SourceConfig,
|
||||||
|
ThinkingObservation,
|
||||||
)
|
)
|
||||||
from polygateway.types import ChatRequest, Effort, LLMResponse, SourceConfig, ThinkingObservation
|
|
||||||
|
|
||||||
|
|
||||||
def _source(**overrides):
|
def _source(**overrides):
|
||||||
@@ -1144,3 +1154,67 @@ async def test_custom_profile_raw_roots_cannot_override_managed_intent(key):
|
|||||||
assert sent == []
|
assert sent == []
|
||||||
finally:
|
finally:
|
||||||
await transport.aclose()
|
await transport.aclose()
|
||||||
|
|
||||||
|
|
||||||
|
class TestRetryAfterNonFinite:
|
||||||
|
"""F1: `Retry-After` 非有限值必须当作"无提示"(计划 §3.6 / §5 批次 G)。
|
||||||
|
|
||||||
|
`float("inf")` 能被 `float()` 成功解析,又能通过既有的 `seconds > 0`——
|
||||||
|
它会一路变成 `retry_after_s=inf`,而 `backoff_delay` 的 `max(delay, retry_after)`
|
||||||
|
取大之后就是一次**永不醒来**的退避 sleep(库不夹 `backoff_max_s`)。
|
||||||
|
"""
|
||||||
|
|
||||||
|
def _translate(self, raw, *, name="qwen_1"):
|
||||||
|
return _translate_429(_source(name=name), "{}", {"retry-after": raw}, {"body_text": "{}"})
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("raw", ["inf", "-inf", "1e999", "Infinity"])
|
||||||
|
def test_non_finite_becomes_no_hint_with_exactly_one_warning(self, raw):
|
||||||
|
messages: list[str] = []
|
||||||
|
sink_id = logger.add(messages.append, level="WARNING")
|
||||||
|
try:
|
||||||
|
exc = self._translate(raw)
|
||||||
|
finally:
|
||||||
|
logger.remove(sink_id)
|
||||||
|
assert exc.retry_after_s is None
|
||||||
|
hits = [m for m in messages if "retry_after_not_finite" in m]
|
||||||
|
assert len(hits) == 1
|
||||||
|
assert "qwen_1" in hits[0] # 告警要能定位到源
|
||||||
|
assert raw not in hits[0] # 但不回显原始头字符串(不拼接、不截断)
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"raw", ["nan", "", " ", "-1", "0", "Wed, 21 Oct 2026 07:28:00 GMT", "soon"]
|
||||||
|
)
|
||||||
|
def test_other_unusable_values_stay_silent(self, raw):
|
||||||
|
"""429 风暴下逐次告警会淹掉真信号: 只有非有限值这一类新增告警。
|
||||||
|
|
||||||
|
`nan` 仍走既有的 `seconds > 0` 恒假值语义,本版**不为它新增分支**。
|
||||||
|
"""
|
||||||
|
messages: list[str] = []
|
||||||
|
sink_id = logger.add(messages.append, level="WARNING")
|
||||||
|
try:
|
||||||
|
exc = self._translate(raw)
|
||||||
|
finally:
|
||||||
|
logger.remove(sink_id)
|
||||||
|
assert exc.retry_after_s is None
|
||||||
|
assert [m for m in messages if "retry_after_not_finite" in m] == []
|
||||||
|
|
||||||
|
def test_finite_positive_still_reaches_backoff(self):
|
||||||
|
"""有限正数一字不改地保留,并照常参与 `max(delay, retry_after)` 取大。"""
|
||||||
|
exc = self._translate("2.5")
|
||||||
|
assert exc.retry_after_s == 2.5
|
||||||
|
policy = RetryPolicy(max_attempts=3, backoff_base_s=0.001, backoff_max_s=0.01)
|
||||||
|
assert backoff_delay(policy, 1, exc, lambda: 0.5) == 2.5
|
||||||
|
# 而非有限值被吃掉之后,退避退回纯指数,不会变成永不醒来的 sleep
|
||||||
|
assert backoff_delay(policy, 1, self._translate("inf"), lambda: 0.5) < 1.0
|
||||||
|
|
||||||
|
def test_source_name_is_a_required_keyword(self):
|
||||||
|
"""私有函数的必填 kw: 漏传即 `TypeError`,不给默认值掩盖调用点漏改。"""
|
||||||
|
with pytest.raises(TypeError):
|
||||||
|
_parse_retry_after("2.5")
|
||||||
|
assert _parse_retry_after("2.5", source_name="qwen_1") == 2.5
|
||||||
|
|
||||||
|
def test_insufficient_quota_is_still_source_dead(self):
|
||||||
|
"""分类判据不受本次改动影响(告警只加在 429 限速那一支)。"""
|
||||||
|
body = json.dumps({"error": {"type": "insufficient_quota"}})
|
||||||
|
exc = _translate_429(_source(), body, {"retry-after": "inf"}, {"body_text": body})
|
||||||
|
assert isinstance(exc, SourceDeadError)
|
||||||
|
|||||||
+136
-1
@@ -76,6 +76,8 @@ class FakeTransport:
|
|||||||
self.script = list(script)
|
self.script = list(script)
|
||||||
self.calls = []
|
self.calls = []
|
||||||
self.efforts = []
|
self.efforts = []
|
||||||
|
# 取消用例的确定性窗口: 进入 hang 分支即置位, 用例据此取消而非 sleep 猜时长
|
||||||
|
self.entered = asyncio.Event()
|
||||||
|
|
||||||
async def complete(self, *, messages, source, stream, overlay, call_id, reasoning_effort):
|
async def complete(self, *, messages, source, stream, overlay, call_id, reasoning_effort):
|
||||||
self.calls.append((source.name, call_id))
|
self.calls.append((source.name, call_id))
|
||||||
@@ -84,10 +86,37 @@ class FakeTransport:
|
|||||||
if isinstance(action, Exception):
|
if isinstance(action, Exception):
|
||||||
raise action
|
raise action
|
||||||
if action == "hang":
|
if action == "hang":
|
||||||
|
self.entered.set()
|
||||||
await asyncio.Event().wait()
|
await asyncio.Event().wait()
|
||||||
return action
|
return action
|
||||||
|
|
||||||
|
|
||||||
|
class HangingGate(InMemoryGate):
|
||||||
|
"""在指定记账写回处永久挂起的门控: 把"取消落在某个 await 上"变成确定性事件。
|
||||||
|
|
||||||
|
只覆盖 `record_success` / `record_failure` 两个写回点, 其余行为沿用真实内存实现。
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, *, hang_on, **kwargs):
|
||||||
|
super().__init__(**kwargs)
|
||||||
|
self._hang_on = hang_on
|
||||||
|
self.entered = asyncio.Event()
|
||||||
|
|
||||||
|
async def _hang(self):
|
||||||
|
self.entered.set()
|
||||||
|
await asyncio.Event().wait()
|
||||||
|
|
||||||
|
async def record_success(self, entry, *, count_attempt=True):
|
||||||
|
if self._hang_on == "success":
|
||||||
|
await self._hang()
|
||||||
|
return await super().record_success(entry, count_attempt=count_attempt)
|
||||||
|
|
||||||
|
async def record_failure(self, entry, reason, force_open):
|
||||||
|
if self._hang_on == "failure":
|
||||||
|
await self._hang()
|
||||||
|
return await super().record_failure(entry, reason, force_open)
|
||||||
|
|
||||||
|
|
||||||
class FakeSleep:
|
class FakeSleep:
|
||||||
"""记录退避时长,立即返回(不真等)。"""
|
"""记录退避时长,立即返回(不真等)。"""
|
||||||
|
|
||||||
@@ -109,6 +138,7 @@ def _harness(
|
|||||||
rng=lambda: 0.0,
|
rng=lambda: 0.0,
|
||||||
selector=None,
|
selector=None,
|
||||||
pacer=None,
|
pacer=None,
|
||||||
|
gate=None,
|
||||||
):
|
):
|
||||||
clock = clock or FakeClock()
|
clock = clock or FakeClock()
|
||||||
limiter = InMemoryLimiter(
|
limiter = InMemoryLimiter(
|
||||||
@@ -118,7 +148,7 @@ def _harness(
|
|||||||
lease_ttl_s=100.0,
|
lease_ttl_s=100.0,
|
||||||
now=clock,
|
now=clock,
|
||||||
)
|
)
|
||||||
gate = InMemoryGate(config=_BREAKER, now=clock)
|
gate = gate if gate is not None else InMemoryGate(config=_BREAKER, now=clock)
|
||||||
transport = FakeTransport(script)
|
transport = FakeTransport(script)
|
||||||
sleep = FakeSleep()
|
sleep = FakeSleep()
|
||||||
mw = RetryMW(
|
mw = RetryMW(
|
||||||
@@ -467,6 +497,111 @@ class TestScopeUnavailable:
|
|||||||
assert resp.content == "ok" and released["done"]
|
assert resp.content == "ok" and released["done"]
|
||||||
|
|
||||||
|
|
||||||
|
class TestCancellationSettlement:
|
||||||
|
"""1.3.6 §6.3 结算矩阵: 取消时 `settle()` 的取值只由"该刻库知道什么"决定。
|
||||||
|
|
||||||
|
取消窗口一律用真实 `asyncio.Event` 钉死(不用 sleep 撞窗口), 否则红绿都不可信。
|
||||||
|
"""
|
||||||
|
|
||||||
|
async def test_cancel_in_flight_keeps_the_reservation(self):
|
||||||
|
"""S3: transport 在途被取消 → 端口已开始、用量未知 → 保留预扣(不凭空退款)。"""
|
||||||
|
mw, limiter, _, transport, *_ = _harness(
|
||||||
|
[_src("a", max_concurrency=1, tpm=1000, est_tokens=400)], ["hang"]
|
||||||
|
)
|
||||||
|
task = asyncio.ensure_future(mw(_REQ))
|
||||||
|
await transport.entered.wait()
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
stats = await limiter.source_stats("a")
|
||||||
|
assert stats.tpm_used == 400 # est 保留, 而非退成 0
|
||||||
|
assert stats.inflight == 0
|
||||||
|
|
||||||
|
async def test_cancel_after_usage_known_keeps_real_usage(self):
|
||||||
|
"""S4: 真实 usage 已算出后被取消 → 结算仍是真实值, 不被 est 覆写。"""
|
||||||
|
clock = FakeClock()
|
||||||
|
gate = HangingGate(hang_on="success", config=_BREAKER, now=clock)
|
||||||
|
mw, limiter, *_ = _harness(
|
||||||
|
[_src("a", max_concurrency=1, tpm=1000, est_tokens=400)],
|
||||||
|
[_ok()],
|
||||||
|
clock=clock,
|
||||||
|
gate=gate,
|
||||||
|
)
|
||||||
|
task = asyncio.ensure_future(mw(_REQ))
|
||||||
|
await gate.entered.wait()
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
assert (await limiter.source_stats("a")).tpm_used == 15 # 10+5 实测
|
||||||
|
|
||||||
|
async def test_cancel_in_dead_failure_branch_keeps_full_refund(self):
|
||||||
|
"""S5-dead: 源已判死时的既有 `0` 不得因取消退化成 est(不得继续占额度)。"""
|
||||||
|
clock = FakeClock()
|
||||||
|
gate = HangingGate(hang_on="failure", config=_BREAKER, now=clock)
|
||||||
|
mw, limiter, *_ = _harness(
|
||||||
|
[_src("a", max_concurrency=1, tpm=1000, est_tokens=400)],
|
||||||
|
[SourceDeadError("401", source_name="a", status_code=401)],
|
||||||
|
clock=clock,
|
||||||
|
gate=gate,
|
||||||
|
)
|
||||||
|
task = asyncio.ensure_future(mw(_REQ))
|
||||||
|
await gate.entered.wait()
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
assert (await limiter.source_stats("a")).tpm_used == 0
|
||||||
|
|
||||||
|
async def test_cancel_in_transient_failure_branch_keeps_the_reservation(self):
|
||||||
|
"""S5-transient: 瞬时失败的结算决定在首个 await 之前定死, 取消拿到同一个 est。"""
|
||||||
|
clock = FakeClock()
|
||||||
|
gate = HangingGate(hang_on="failure", config=_BREAKER, now=clock)
|
||||||
|
mw, limiter, *_ = _harness(
|
||||||
|
[_src("a", max_concurrency=1, tpm=1000, est_tokens=400)],
|
||||||
|
[TransientError("boom", source_name="a", status_code=500)],
|
||||||
|
clock=clock,
|
||||||
|
gate=gate,
|
||||||
|
)
|
||||||
|
task = asyncio.ensure_future(mw(_REQ))
|
||||||
|
await gate.entered.wait()
|
||||||
|
task.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
assert (await limiter.source_stats("a")).tpm_used == 400
|
||||||
|
|
||||||
|
async def test_unclassified_exception_still_refunds_in_full(self):
|
||||||
|
"""S8 防越界: 未分类异常(无 except 接住)仍逐字走 1.3.5 的全额退还。"""
|
||||||
|
mw, limiter, *_ = _harness(
|
||||||
|
[_src("a", max_concurrency=1, tpm=1000, est_tokens=400)], [RuntimeError("boom")]
|
||||||
|
)
|
||||||
|
with pytest.raises(RuntimeError):
|
||||||
|
await mw(_REQ)
|
||||||
|
stats = await limiter.source_stats("a")
|
||||||
|
assert stats.tpm_used == 0 and stats.inflight == 0
|
||||||
|
|
||||||
|
async def test_real_zero_usage_success_settles_zero(self):
|
||||||
|
"""防越界: 真实 usage 恰为 0 是**已知事实**, 不得被当成"未知"改按 est 结算。"""
|
||||||
|
zero = dataclasses.replace(_ok(), prompt_tokens=0, completion_tokens=0)
|
||||||
|
mw, limiter, *_ = _harness([_src("a", tpm=1000, est_tokens=400)], [zero])
|
||||||
|
await mw(_REQ)
|
||||||
|
assert (await limiter.source_stats("a")).tpm_used == 0
|
||||||
|
|
||||||
|
async def test_circuit_open_rejection_settles_zero_end_to_end(self):
|
||||||
|
"""S1 端到端: 开路拒绝的 pick 预扣后按 0 结算, 不给 tpm_used 增加任何量。"""
|
||||||
|
clock = FakeClock()
|
||||||
|
script = [TransientError(str(i)) for i in range(9)]
|
||||||
|
mw, limiter, *_ = _harness(
|
||||||
|
[_src("a", max_concurrency=1, tpm=10000, est_tokens=400)],
|
||||||
|
script,
|
||||||
|
clock=clock,
|
||||||
|
max_attempts=99,
|
||||||
|
)
|
||||||
|
# 3 次瞬时失败后 a 开路 → 第 4 次 pick 被拒绝
|
||||||
|
with pytest.raises(CircuitOpenError):
|
||||||
|
await mw(_REQ)
|
||||||
|
# 三次瞬时失败各保留 est = 1200; 开路那次 pick 若漏了 settle(0) 会再 +400
|
||||||
|
assert (await limiter.source_stats("a")).tpm_used == 1200
|
||||||
|
|
||||||
|
|
||||||
class TestCancellation:
|
class TestCancellation:
|
||||||
async def test_cancel_mid_flight_releases_permit(self):
|
async def test_cancel_mid_flight_releases_permit(self):
|
||||||
mw, limiter, _, _, _, _ = _harness([_src("a", max_concurrency=1)], ["hang"])
|
mw, limiter, _, _, _, _ = _harness([_src("a", max_concurrency=1)], ["hang"])
|
||||||
|
|||||||
Reference in New Issue
Block a user