diff --git a/src/polygateway/backends/redis/breaker.py b/src/polygateway/backends/redis/breaker.py index 3fe0de8..dd196bd 100644 --- a/src/polygateway/backends/redis/breaker.py +++ b/src/polygateway/backends/redis/breaker.py @@ -367,7 +367,9 @@ class RedisGate: keys=[self._key(source_name)], args=[owner, self._probe_ttl_ms] ) except RedisError as exc: - raise GovernanceBackendError(f"熔断后端 try_enter 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"熔断后端 try_enter 失败: {exc}", scope=self._scope + ) from exc return self._decision(source_name, result) async def record_success( @@ -385,7 +387,9 @@ class RedisGate: try: result = await self._success_lua(keys=[self._key(entry.source_name)], args=args) except RedisError as exc: - raise GovernanceBackendError(f"熔断后端 record_success 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"熔断后端 record_success 失败: {exc}", scope=self._scope + ) from exc return self._update(result) async def record_failure( @@ -407,7 +411,9 @@ class RedisGate: try: result = await self._failure_lua(keys=[self._key(entry.source_name)], args=args) except RedisError as exc: - raise GovernanceBackendError(f"熔断后端 record_failure 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"熔断后端 record_failure 失败: {exc}", scope=self._scope + ) from exc return self._update(result) async def release_probe(self, entry: GateDecision) -> GateUpdate: @@ -419,7 +425,9 @@ class RedisGate: keys=[self._key(entry.source_name)], args=[entry.epoch, entry.probe_owner] ) except RedisError as exc: - raise GovernanceBackendError(f"熔断后端 release_probe 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"熔断后端 release_probe 失败: {exc}", scope=self._scope + ) from exc return self._update(result) async def retry_after_s(self, sources: tuple[str, ...]) -> float: @@ -429,7 +437,9 @@ class RedisGate: try: result = await self._retry_after_lua(keys=[self._key(s) for s in sources]) except RedisError as exc: - raise GovernanceBackendError(f"熔断后端 retry_after_s 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"熔断后端 retry_after_s 失败: {exc}", scope=self._scope + ) from exc return int(result) / 1000.0 async def aclose(self) -> None: diff --git a/src/polygateway/backends/redis/limiter.py b/src/polygateway/backends/redis/limiter.py index 011e925..cb5d516 100644 --- a/src/polygateway/backends/redis/limiter.py +++ b/src/polygateway/backends/redis/limiter.py @@ -247,7 +247,9 @@ class RedisLimiter: ], ) except RedisError as exc: - raise GovernanceBackendError(f"限流后端 try_acquire 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"限流后端 try_acquire 失败: {exc}", scope=self._scope + ) from exc if ok != 1: return None return _RedisPermit(self, source_key, lease_id, est_tokens, window) @@ -265,7 +267,9 @@ class RedisLimiter: try: await self._release_lua(keys=[gl, sl], args=[lease_id]) except RedisError as exc: - raise GovernanceBackendError(f"限流后端 release 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"限流后端 release 失败: {exc}", scope=self._scope + ) from exc async def _settle_tpm(self, source_key: str, delta: int, window: int) -> None: wk = self._window_keys(source_key, window) @@ -283,7 +287,9 @@ class RedisLimiter: wk = self._window_keys(source_key, window) res = await self._stats_lua(keys=[sl, wk["s_rpm"], wk["s_tpm"]]) except RedisError as exc: - raise GovernanceBackendError(f"限流后端 source_stats 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"限流后端 source_stats 失败: {exc}", scope=self._scope + ) from exc return SourceStats( inflight=int(res[0]), rpm_used=max(0, int(res[1])), @@ -295,14 +301,18 @@ class RedisLimiter: try: await self._progress_mark_lua(keys=[self._progress_key()], args=[_PROGRESS_TTL_S]) except RedisError as exc: - raise GovernanceBackendError(f"限流后端 mark_progress 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"限流后端 mark_progress 失败: {exc}", scope=self._scope + ) from exc async def progress_age_s(self) -> float: """距上次全局成功的秒数;仅键缺失(-1)= 从未进展 → inf(CHS limiter.py:208)。""" try: res = await self._progress_age_lua(keys=[self._progress_key()]) except RedisError as exc: - raise GovernanceBackendError(f"限流后端 progress_age_s 失败: {exc}", scope=self._scope) from exc + raise GovernanceBackendError( + f"限流后端 progress_age_s 失败: {exc}", scope=self._scope + ) from exc return float("inf") if int(res) == -1 else int(res) / 1000.0 async def aclose(self) -> None: diff --git a/tests/unit/test_backpressure.py b/tests/unit/test_backpressure.py index 774e904..a8d1d1a 100644 --- a/tests/unit/test_backpressure.py +++ b/tests/unit/test_backpressure.py @@ -321,9 +321,7 @@ class TestStallBudget: async def advance(_n): clock.advance(_STALL) - mw = _mw( - [src], limiter, [], clock=clock, sleep=BoundedSleep(advance), transport=transport - ) + mw = _mw([src], limiter, [], clock=clock, sleep=BoundedSleep(advance), transport=transport) with pytest.raises(AllSourcesExhausted) as ei: await mw(_REQ) assert ei.value.reason == "stalled" # 不是 retry_exhausted: 429 确实没烧重试预算 @@ -541,9 +539,7 @@ class TestUnknownSourceIsAssemblyDefect: """ src = make_source("s1") # 限流后端的源名单与治理循环拿到的源对不上 = 装配缺陷 - limiter = InMemoryLimiter( - scope="llm", sources={"other": src}, global_limits=_NO_GLOBAL - ) + limiter = InMemoryLimiter(scope="llm", sources={"other": src}, global_limits=_NO_GLOBAL) gate = QuotaGate(limiter, scope="llm") with pytest.raises(SourceNotConfiguredError) as ei: await getattr(gate, method)(src)