From 0b8460b6d5cf1e70157988a8d19b83cbba81e240 Mon Sep 17 00:00:00 2001 From: iomgaa Date: Mon, 20 Jul 2026 22:01:26 -0400 Subject: [PATCH] fix: address independent verification findings Classify empty completions as transient per human ruling (fixes flaky real-gateway smoke and prevents caching empty responses), rename the factory injection parameter gate to breaker per the frozen design, rewrite the probe-entry cleanup without except BaseException, declare python-dotenv explicitly, add a mid-backoff cancellation test, and record all implementation errata in the design and architecture docs. --- pyproject.toml | 3 ++ research-wiki/ARCHITECTURE.md | 6 ++-- .../designs/2026-07-20-m1-core-design.md | 3 +- src/polygateway/backends/memory/limiter.py | 3 +- src/polygateway/client.py | 12 +++---- src/polygateway/middleware/retry.py | 8 +++-- src/polygateway/transports/openai_compat.py | 16 +++++++++ tests/integration/test_governance_stack.py | 2 +- tests/integration/test_redis_cache.py | 2 +- tests/unit/test_client.py | 2 +- tests/unit/test_openai_compat.py | 27 +++++++++++++++ tests/unit/test_retry.py | 33 +++++++++++++++++++ 12 files changed, 101 insertions(+), 16 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index c90847d..cf5a079 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -11,6 +11,9 @@ dependencies = [ "httpx>=0.27", "pydantic>=2.8", "pydantic-settings>=2.4", + # 显式声明(2026-07-20 人类确认): config.py 的多源动态键族直接使用 dotenv_values, + # 不能只依赖 pydantic-settings 的传递安装 + "python-dotenv>=1.0", "loguru>=0.7", ] diff --git a/research-wiki/ARCHITECTURE.md b/research-wiki/ARCHITECTURE.md index 2b186f7..60f95c7 100644 --- a/research-wiki/ARCHITECTURE.md +++ b/research-wiki/ARCHITECTURE.md @@ -362,6 +362,7 @@ flowchart TB | HTTP 429(body 无 insufficient_quota)、500/502/503/504 | `TransientError`(携 `Retry-After` 解析值,仅支持秒数形态) | | HTTP 429 + body 含 insufficient_quota、401、403 | `SourceDeadError` | | HTTP 400 | `RequestRejectedError` | +| **空补全**: 200 且流程完整([DONE]/usage 正常)但 content 为空(2026-07-20 M1 验证发现,人类裁决) | `TransientError`(服务抖动,重试/换源;绝不缓存空响应) | | 解析层失败(结构化输出/OCR ZIP) | `ResultInvalidError` | ### 6.3 "坏结果 ≠ 坏服务"(ResultInvalidError 语义,继承 CHSAnalyzer) @@ -396,7 +397,7 @@ flowchart TB - `InMemoryLimiter`: 同一契约的进程内实现(semaphore + 滑动窗口计数);单进程场景下语义等价。 - **配额满行为可配**: `wait`(等待,配 stall 判定——本地等待超窗 + 全局无进展超窗双条件才判卡死)或 `fail-fast`(立即抛)。 - 全局活性信号: `mark_progress()`/`progress_age_s()`("最近一次出餐"时刻)供背压 stall 判定,移植 `CHSAnalyzer limiter.py:193`。 -- **契约补强(2026-07-20,CHS 迁移缺口 G6)**: `settle()`/`release()` 幂等(重复调用无副作用);装配期守卫——`timeout_s ≤ permit 租约 TTL`(防租约先于请求过期)、`stall_window ≥ 最慢源 TTFT 上限`(防误判卡死),违反直接报错拒绝装配。 +- **契约补强(2026-07-20,CHS 迁移缺口 G6)**: `settle()`/`release()` 幂等(重复调用无副作用);装配期守卫——`timeout_s ≤ permit 租约 TTL`(防租约先于请求过期)、`stall_window ≥ 最慢源 TTFT 上限`(防误判卡死),违反直接报错拒绝装配。降级方向细化(2026-07-20 M1): "报错不放行"适用于**准入侧**(try_acquire/try_enter 等);已成功调用后的 settle/release 释放侧失败降级 warning——释放失败不构成放行,且不得掩盖主异常与取消。 ### 7.4 熔断 @@ -469,6 +470,7 @@ src/polygateway/ ├── errors.py # §6 错误四分类 ├── ports.py # 全部 Protocol(§4 各端口) ├── client.py # GatewayClient + from_env()/from_settings() 装配工厂 + gather_bounded +├── config.py # GatewaySettings: 多源/韧性/装配键族聚合与装配守卫(M1 增补) ├── middleware/ # retry.py / ratelimit.py / breaker.py / cache.py / telemetry.py / structured.py ├── transports/ # openai_compat.py / openai_sdk.py / monkey_ocr.py ├── providers.py # D11 provider 注册表 @@ -485,7 +487,7 @@ src/polygateway/ ## 9. 配置面 -- **载体**: `pydantic-settings` + `.env`(工程配置);缺失关键配置直接报错,严禁硬编码默认值兜底(三项目共同铁律)。 +- **载体**: `.env` + 环境变量(工程配置);缺失关键配置直接报错,严禁硬编码默认值兜底(三项目共同铁律)。**实现勘误(2026-07-20 M1,人类确认)**: 多源 `{SCOPE}__{PROVIDER}__{N}__{FIELD}` 是动态键族,pydantic-settings 的静态字段模型无法表达,故 `GatewaySettings` 为 frozen dataclass + python-dotenv(显式核心依赖)读取,fail-loud 校验语义与 pydantic-settings 一致。 - **多源命名**: `{SCOPE}__{PROVIDER}__{N}__{FIELD}`(如 `LLM__QWEN__1__API_KEY`、`OCR__MONKEY__1__BASE_URL`),聚合为 `list[SourceConfig]`;SCOPE 支持逻辑角色前缀(§7.7)。 - **韧性参数键名**沿用三项目习惯(`LLM_TIMEOUT` / `LLM_MAX_RETRIES` / `LLM_RETRY_BASE_DELAY` / `LLM_RETRY_MAX_DELAY` / `LLM_CIRCUIT_BREAKER_THRESHOLD` / `LLM_CIRCUIT_BREAKER_COOLDOWN` / `LLM_TTFT_TIMEOUT` / `LLM_INTER_TOKEN_TIMEOUT`),降低三项目迁移改名成本。 - **per-scope 韧性配置(2026-07-20,CHS 迁移缺口 G4)**: 韧性参数支持按 scope 覆盖——`{SCOPE}__RETRY__MAX_ATTEMPTS` / `{SCOPE}__BREAKER__FAIL_THRESHOLD` / `{SCOPE}__BREAKER__COOLDOWN_S` / `{SCOPE}__BACKPRESSURE__STALL_WINDOW_S` / `{SCOPE}__SELECTOR` / `{SCOPE}__GLOBAL__MAX_CONCURRENCY|RPM|TPM`(CHS 现状: VLM 与 OCR 两 scope 参数各异)。平铺键(`LLM_*`)是单 scope 场景的简写;两者并存时 scope 键优先。 diff --git a/research-wiki/designs/2026-07-20-m1-core-design.md b/research-wiki/designs/2026-07-20-m1-core-design.md index 7dba1b3..9fbac94 100644 --- a/research-wiki/designs/2026-07-20-m1-core-design.md +++ b/research-wiki/designs/2026-07-20-m1-core-design.md @@ -1,6 +1,7 @@ # M1 核心里程碑设计:公共签名冻结与治理栈落地 -> **状态**: Claude 自审 ✅ → 独立审查 ✅(全新上下文 subagent 顶替 Codex——本机 codex CLI 损坏;6 Issues + 4 建议已全部核验修订)→ **待人类门**。**依据**: ARCHITECTURE.md D1–D14 / ROADMAP §2 / migrations/ 三份;签名证据来自 2026-07-20 对 reference/ 三项目的逐字提取(文中 `文件:行号` 均指 `reference/` 下路径)。 +> **状态**: Claude 自审 ✅ → 独立审查 ✅ → 人类门 ✅(2026-07-20)→ 实现完成 → 独立验证 ✅(verifier 发现 1 Critical + 3 Important 均已处理)。 +> **实现期勘误(2026-07-20,均经人类确认)**: ① **空补全分类**——服务 200 且流程完整但 content 为空(MiniMax 间歇形态)→ `TransientError("empty_completion")`,退避重试/换源,绝不缓存(承 CHS 零内容归瞬时先例;验证 e2e flaky 的根因修复);② `GatewaySettings` 实现为 frozen dataclass + python-dotenv(多源动态键族 pydantic-settings 无法建模;python-dotenv 显式入核心依赖),fail-loud 语义不变;③ `minimax` 基线 profile 入 DEFAULT_PROFILES(首个 register 之外的新增条目,D11 实战);④ §5 细则 2 升级策略实现为"全部源都声明 supports_native_schema 才启用 escalation"(混合支持的 scope 注入 response_format 会击穿不支持的源)。**依据**: ARCHITECTURE.md D1–D14 / ROADMAP §2 / migrations/ 三份;签名证据来自 2026-07-20 对 reference/ 三项目的逐字提取(文中 `文件:行号` 均指 `reference/` 下路径)。 > **范围**: M1 全部交付物(ROADMAP §2 七步),含 2026-07-20 人类拍板:多源完整行为进 M1;GovDoc 与 Video-Tree 双冒烟;Redis 集成测试用实验室远程实例。 > **本文档冻结的内容**: `types.py` / `errors.py` / `ports.py` 公共签名、配置键名、五个开放设计轴的取舍。与 ARCHITECTURE.md 冲突处在 §13 列为反哺修订,经人类批准后先改 ARCHITECTURE 再实施。 diff --git a/src/polygateway/backends/memory/limiter.py b/src/polygateway/backends/memory/limiter.py index 278bdd8..5493bc5 100644 --- a/src/polygateway/backends/memory/limiter.py +++ b/src/polygateway/backends/memory/limiter.py @@ -1,7 +1,8 @@ """进程内限流后端: 与 Redis 版同一契约的六道闸实现(D3 双后端)。 语义蓝本 CHS `app/coordination/limiter.py`: 并发 = 带 TTL 的租约(持有者 -死亡后过期回收);RPM/TPM = 分钟滑动窗口计数;TPM 入场按 est 预扣,settle +死亡后过期回收);RPM/TPM = 分钟**固定窗口**计数(`int(now/60)`,与 CHS +Lua 同款口径);TPM 入场按 est 预扣,settle 按实际结算多退少补且退款落 acquire 时的窗口。检查-占用在单次同步段内完成 (无 await 穿插),单事件循环下天然原子;本实现不跨进程,是单进程部署的 正确答案(多 worker 用 M2 Redis 后端)。 diff --git a/src/polygateway/client.py b/src/polygateway/client.py index ce9fb3e..7ac5665 100644 --- a/src/polygateway/client.py +++ b/src/polygateway/client.py @@ -62,7 +62,7 @@ class GatewayClient: sources: list[SourceConfig], selector: SourceSelector, limiter: RateLimiter, - gate: ProviderGate, + breaker: ProviderGate, transport: Transport, retry: RetryPolicy, backpressure: BackpressurePolicy, @@ -84,7 +84,7 @@ class GatewayClient: sources=sources, selector=selector, limiter=limiter, - gate=gate, + gate=breaker, transport=transport, retry=retry, backpressure=backpressure, @@ -183,7 +183,7 @@ class GatewayClient: settings: GatewaySettings, *, limiter: RateLimiter | None = None, - gate: ProviderGate | None = None, + breaker: ProviderGate | None = None, cache: CacheBackend | None = None, telemetry: TelemetryRecorder | None = None, registry: Mapping[str, ProviderProfile] | None = None, @@ -203,7 +203,7 @@ class GatewayClient: global_limits=settings.global_limits, lease_ttl_s=settings.lease_ttl_s, ), - gate=gate or InMemoryGate(config=settings.breaker), + breaker=breaker or InMemoryGate(config=settings.breaker), transport=OpenAICompatTransport(registry=registry), retry=settings.retry, backpressure=settings.backpressure, @@ -223,7 +223,7 @@ class GatewayClient: scope: str = "LLM", *, limiter: RateLimiter | None = None, - gate: ProviderGate | None = None, + breaker: ProviderGate | None = None, cache: CacheBackend | None = None, telemetry: TelemetryRecorder | None = None, registry: Mapping[str, ProviderProfile] | None = None, @@ -233,7 +233,7 @@ class GatewayClient: return cls.from_settings( GatewaySettings.from_env(scope, env=env), limiter=limiter, - gate=gate, + breaker=breaker, cache=cache, telemetry=telemetry, registry=registry, diff --git a/src/polygateway/middleware/retry.py b/src/polygateway/middleware/retry.py index e691c4b..6a681ce 100644 --- a/src/polygateway/middleware/retry.py +++ b/src/polygateway/middleware/retry.py @@ -154,11 +154,13 @@ class RetryMW: if permit is None: reasons.setdefault(cand.name, "rate_limited") continue + entry = None try: entry = await self._breaker.try_enter(cand, uuid.uuid4().hex) - except BaseException: - await self._settle_and_release(permit, 0) - raise + finally: + # try_enter 未归还 entry(异常/取消)→ 释放已占 permit,不吞任何异常 + if entry is None: + await self._settle_and_release(permit, 0) if entry.allowed: return (cand, permit, entry), gate_rejections gate_rejections += 1 diff --git a/src/polygateway/transports/openai_compat.py b/src/polygateway/transports/openai_compat.py index a50c453..1dc56ac 100644 --- a/src/polygateway/transports/openai_compat.py +++ b/src/polygateway/transports/openai_compat.py @@ -252,6 +252,7 @@ class OpenAICompatTransport: (content_parts if is_content else thinking_parts).append(text) salvaged = self._check_done(sink, content_parts, thinking_parts, source) content, thinking = self._finalize_text(content_parts, thinking_parts, profile) + self._reject_empty_completion(content, source) prompt, completion, usage_source = _resolve_usage(sink.get("usage") or {}, source) if salvaged: usage_source = "estimated" # 打捞路径强制 estimated(设计 §6) @@ -283,6 +284,20 @@ class OpenAICompatTransport: raise TransientError(f"{source.name} SSE missing_done: 截断且无 [DONE]", **ctx) return True + def _reject_empty_completion(self, content: str, source: SourceConfig) -> None: + """空补全 → 瞬时错误(2026-07-20 人类裁决,M1 验证发现)。 + + 服务 200 且流程完整([DONE]/usage 正常)但 content 为空——MiniMax 等 + 网关的间歇异常形态。视为服务抖动: 退避重试/换源,**绝不缓存空响应**; + 承 CHS "VLM 零 content"归瞬时的先例(invokers.py:309)。 + """ + if not content.strip(): + raise TransientError( + f"{source.name} 空补全(empty_completion): 流程完整但零内容", + source_name=source.name, + operation="chat", + ) + def _finalize_text( self, content_parts: list[str], thinking_parts: list[str], profile: ProviderProfile ) -> tuple[str, str]: @@ -321,6 +336,7 @@ class OpenAICompatTransport: content, thinking = self._finalize_text( [message.get("content") or ""], [message.get("reasoning_content") or ""], profile ) + self._reject_empty_completion(content, source) prompt, completion, usage_source = _resolve_usage(body.get("usage") or {}, source) return TransportResult( content=content, diff --git a/tests/integration/test_governance_stack.py b/tests/integration/test_governance_stack.py index b93c700..86eb237 100644 --- a/tests/integration/test_governance_stack.py +++ b/tests/integration/test_governance_stack.py @@ -65,7 +65,7 @@ def _full_client(handler, *, clock=None, telemetry=None, cache=None): global_limits=GlobalLimits(0, 0, 0), now=clock, ), - gate=InMemoryGate(config=_BREAKER, now=clock), + breaker=InMemoryGate(config=_BREAKER, now=clock), transport=OpenAICompatTransport( client_factory=lambda s: httpx.AsyncClient(transport=httpx.MockTransport(handler)) ), diff --git a/tests/integration/test_redis_cache.py b/tests/integration/test_redis_cache.py index d559be8..8ebf40a 100644 --- a/tests/integration/test_redis_cache.py +++ b/tests/integration/test_redis_cache.py @@ -155,7 +155,7 @@ class TestClientWithRealRedis: limiter=InMemoryLimiter( scope="llm", sources={src.name: src}, global_limits=GlobalLimits(0, 0, 0) ), - gate=InMemoryGate(config=BreakerConfig(5, 60.0, 120.0)), + breaker=InMemoryGate(config=BreakerConfig(5, 60.0, 120.0)), transport=OpenAICompatTransport( client_factory=lambda s: httpx.AsyncClient( transport=httpx.MockTransport(lambda req: _sse_response()) diff --git a/tests/unit/test_client.py b/tests/unit/test_client.py index 7203871..080112d 100644 --- a/tests/unit/test_client.py +++ b/tests/unit/test_client.py @@ -81,7 +81,7 @@ def _client(sources=None, handler=None, *, limiter=None, quota_full="wait", **ov sources={s.name: s for s in sources}, global_limits=GlobalLimits(0, 0, 0), ), - "gate": InMemoryGate(config=BreakerConfig(5, 60.0, 120.0)), + "breaker": InMemoryGate(config=BreakerConfig(5, 60.0, 120.0)), "transport": transport, "retry": RetryPolicy(3, 2.0, 30.0), "backpressure": BackpressurePolicy(300.0, 0.01), diff --git a/tests/unit/test_openai_compat.py b/tests/unit/test_openai_compat.py index 3aee2d7..0f16eda 100644 --- a/tests/unit/test_openai_compat.py +++ b/tests/unit/test_openai_compat.py @@ -161,6 +161,33 @@ class TestMissingDoneSemantics: await _complete(_transport_for(handler), _source(missing_done="salvage")) +class TestEmptyCompletion: + """空补全 → TransientError(2026-07-20 人类裁决;MiniMax 间歇形态,绝不缓存)。""" + + async def test_stream_zero_content_with_done_is_transient(self): + def handler(request): + return _sse_stream(_chunk(reasoning="only thinking"), _chunk(usage=_USAGE)) + + with pytest.raises(TransientError, match="empty_completion"): + await _complete(_transport_for(handler), _source()) + + async def test_non_stream_empty_content_is_transient(self): + def handler(request): + return httpx.Response( + 200, json={"choices": [{"message": {"content": ""}}], "usage": _USAGE} + ) + + with pytest.raises(TransientError, match="empty_completion"): + await _complete(_transport_for(handler), _source(), stream=False) + + async def test_think_only_content_after_strip_is_transient(self): + def handler(request): + return _sse_stream(_chunk(content="hmm"), _chunk(usage=_USAGE)) + + with pytest.raises(TransientError, match="empty_completion"): + await _complete(_transport_for(handler), _source()) + + class TestNonStreamFastPath: async def test_non_stream_parses_message(self): def handler(request): diff --git a/tests/unit/test_retry.py b/tests/unit/test_retry.py index 346c5cb..13a3293 100644 --- a/tests/unit/test_retry.py +++ b/tests/unit/test_retry.py @@ -298,6 +298,39 @@ class TestCancellation: await task assert (await limiter.source_stats("a")).inflight == 0 # finally 释放 + async def test_cancel_mid_backoff_propagates_with_no_held_permit(self): + """退避 sleep 中取消: CancelledError 穿透,且 permit 早已在 finally 释放。""" + clock = FakeClock() + src = _src("a", max_concurrency=1) + limiter = InMemoryLimiter( + scope="llm", + sources={"a": src}, + global_limits=_NO_GLOBAL, + lease_ttl_s=100.0, + now=clock, + ) + mw = RetryMW( + scope="llm", + sources=[src], + selector=RoundRobinSelector(), + limiter=limiter, + gate=InMemoryGate(config=_BREAKER, now=clock), + transport=FakeTransport([TransientError("x"), _ok()]), + retry=RetryPolicy(max_attempts=3, backoff_base_s=30.0, backoff_max_s=60.0), + backpressure=BackpressurePolicy(300.0, 0.01), + cooldown_memo=SourceCooldownMemo(now=clock), + emitter=None, + now=clock, + sleep=asyncio.sleep, + rng=lambda: 0.5, + ) + task = asyncio.ensure_future(mw(_REQ)) + await asyncio.sleep(0.05) # 第一次失败后进入 30s 真实退避 + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + assert (await limiter.source_stats("a")).inflight == 0 # 退避期不占并发槽 + async def test_cancel_probe_releases_probe_lease(self): clock = FakeClock() mw, _, gate, _, _, _ = _harness(