"""选源策略与源冷却备忘(M1 设计 §2.3;蓝本 CHS app/providers/selector.py 与 governance.py:107)。 选源是端口(SourceSelector),两个首发实现逐字移植 CHS;冷却备忘是 RetryMW 的进程本地状态——熔断开路的源在本地记冷却截止,选源时跳过, 避免每轮白烧 RPM 去探测已知开路的源。 """ from __future__ import annotations import random import time from typing import TYPE_CHECKING if TYPE_CHECKING: from collections.abc import Callable from polygateway.types import SourceConfig, SourceStats _EWMA_ALPHA = 0.2 # 健康 EWMA 步长: 约 10 次成功从谷底爬回 0.9(天然 slow-start) _SCORE_FLOOR = 0.05 # 探索地板: 塌陷源保有微量被选概率,恢复靠真实成功自证 class RoundRobinSelector: """轮转起点后移(CHS selector.py:20 同款);单 client 内游标推进。""" def __init__(self) -> None: self._n = 0 def order( self, sources: list[SourceConfig], stats: dict[str, SourceStats] ) -> list[SourceConfig]: if not sources: return [] k = self._n % len(sources) self._n += 1 return sources[k:] + sources[:k] class LeastInflightSelector: """最少在途优先(CHS selector.py:36 同款);排序稳定,平局保持配置序。""" def order( self, sources: list[SourceConfig], stats: dict[str, SourceStats] ) -> list[SourceConfig]: return sorted(sources, key=lambda s: stats[s.name].inflight if s.name in stats else 0) class HealthAwareSelector: """健康感知选源(M2.5 设计 §3.2;蓝本 Envoy least-request + gRPC WRR)。 score = max(ewma_success, 地板) / (1 + inflight)。头名经 P2C(随机取 两源比分,高者先)引入探索;其余按分数降序。健康态为进程本地(业界 共识: Envoy/Finagle/gRPC 全本地),属 client 实例,不违反纯 asyncio 中立。 已知取舍(设计 §3.2): 仅头名随机化,多 worker 溢出会集中到同一次优源。 """ def __init__(self, *, rng: Callable[[], float] = random.random) -> None: self._rng = rng self._ewma: dict[str, float] = {} def record_outcome(self, source_name: str, ok: bool) -> None: """尝试结果喂数(OutcomeAwareSelector 端口);初始 1.0 乐观起步。""" prev = self._ewma.get(source_name, 1.0) self._ewma[source_name] = prev + _EWMA_ALPHA * ((1.0 if ok else 0.0) - prev) def health(self, source_name: str) -> float: """EWMA 裸值(OutcomeAwareSelector 端口): RetryMW 降权门槛用。""" return self._ewma.get(source_name, 1.0) def _score(self, name: str, stats: dict[str, SourceStats]) -> float: ewma = max(self._ewma.get(name, 1.0), _SCORE_FLOOR) inflight = stats[name].inflight if name in stats else 0 return ewma / (1.0 + inflight) def order( self, sources: list[SourceConfig], stats: dict[str, SourceStats] ) -> list[SourceConfig]: if len(sources) < 2: return list(sources) ranked = sorted(sources, key=lambda s: self._score(s.name, stats), reverse=True) # P2C: 随机取两源比分,胜者提为头名(平分取采样序首位) i = int(self._rng() * len(sources)) % len(sources) j = int(self._rng() * len(sources)) % len(sources) a, b = sources[i], sources[j] head = a if self._score(a.name, stats) >= self._score(b.name, stats) else b return [head] + [s for s in ranked if s.name != head.name] class AdaptivePacer: """AIMD 自适应并发(M2.5 设计 §3.35;Netflix concurrency-limits 损失型)。 429 是网关的"降速"信号: 乘性削减该源并发上限(×0.7),真实成功加性 增长(+1/limit),上限收敛到网关可持续水位;超限调用在 RetryMW 的 quota-wait 轮询里排队而非烧重试预算。进程本地,属 client 实例。 """ _INITIAL = 8.0 _CUT = 0.5 # M2.5 迭代4: 0.7→0.5,残漏 5.8% 的 429 证明在临界点上方震荡 _FLOOR = 1.0 def __init__(self, *, ceiling: float) -> None: if ceiling < self._FLOOR: raise ValueError("ceiling 不得小于下限 1") self._ceiling = ceiling self._limit: dict[str, float] = {} self._inflight: dict[str, int] = {} def limit(self, source_name: str) -> float: return self._limit.get(source_name, min(self._INITIAL, self._ceiling)) def on_backpressure(self, source_name: str) -> None: self._limit[source_name] = max(self._FLOOR, self.limit(source_name) * self._CUT) def on_success(self, source_name: str) -> None: cur = self.limit(source_name) self._limit[source_name] = min(self._ceiling, cur + 1.0 / cur) def admit(self, source_name: str) -> bool: return self._inflight.get(source_name, 0) < self.limit(source_name) def enter(self, source_name: str) -> None: self._inflight[source_name] = self._inflight.get(source_name, 0) + 1 def leave(self, source_name: str) -> None: self._inflight[source_name] = max(0, self._inflight.get(source_name, 0) - 1) class SourceCooldownMemo: """进程本地的源冷却备忘(CHS governance.py:107 同款)。 只记"冷却截止时刻";set_until 取更晚者,防止较早的提示回退已有备忘。 """ def __init__(self, now: Callable[[], float] = time.monotonic) -> None: self._now = now self._until: dict[str, float] = {} def set_until(self, source_name: str, until: float) -> None: self._until[source_name] = max(self._until.get(source_name, 0.0), until) def active(self, source_name: str) -> bool: return self._until.get(source_name, 0.0) > self._now() def remaining(self, source_name: str) -> float: return max(0.0, self._until.get(source_name, 0.0) - self._now())