31 KiB
M2 分布式里程碑实现计划
目标: 按已批准设计
designs/2026-07-20-m2-distributed-design.md交付 Redis 限流/熔断后端、背压 stall、Postgres 遥测、pricing、Embedding 客户端与真实数据压测 harness。 方案概述: CHS 六道闸与熔断 Lua 逐字移植进backends/redis/,经 M1 预留的契约 fixture 参数化接入(时间语义用例走真实等待变体);stall 与记账侧降级改middleware/retry.py;Postgres/pricing/Embedding 为纯新增(只增不改冻结面);harness 落tools/soak/。 涉及技术: redis.asyncio + Lua、asyncpg、httpx、multiprocessing(harness)。 执行环境: conda envPolyGateway;真实 Redis =.envREDIS_URL(db3 专用);真实 Postgres =.envPGW_TELEMETRY_PG_DSN(polygateway 专用库,该实例 app/chs_prod/mimiciv 等在用库严禁触碰)。 纪律: 每任务先写失败测试再实现(测试结果门);reference/ 只读;不做计划外重构。
文件结构(创建 C / 修改 M)
| 文件 | C/M | 职责 |
|---|---|---|
src/polygateway/backends/redis/__init__.py |
C | 导出 RedisLimiter/RedisGate |
src/polygateway/backends/redis/limiter.py |
C | 六道闸 Lua + RedisLimiter/_RedisPermit |
src/polygateway/backends/redis/breaker.py |
C | 熔断五操作 Lua + RedisGate |
src/polygateway/middleware/retry.py |
M | stall 双条件 + poll jitter + 记账侧降级捕获 + _backoff_delay 提取为模块级纯函数 |
src/polygateway/middleware/ratelimit.py |
M | QuotaGate 增 progress_age_s() 透传(stall 判定消费) |
src/polygateway/telemetry/postgres.py |
C | PostgresRecorder(18 列,两级降级) |
src/polygateway/pricing.py |
C | ModelPrice/PricingTable |
src/polygateway/middleware/telemetry.py |
M | TelemetryEmitter 增 pricing 注入,cost 换算 |
src/polygateway/types.py |
M | 只增:EmbeddingResponse、EmbeddingTransportResult |
src/polygateway/ports.py |
M | 只增:EmbeddingTransport Protocol |
src/polygateway/transports/openai_compat.py |
M | 增 embed()(POST /embeddings,复用错误翻译) |
src/polygateway/embedding.py |
C | EmbeddingClient(精简治理循环 + from_env/from_settings) |
src/polygateway/config.py |
M | 解禁 redis 后端键、telemetry 白名单扩 postgres、新键 PG_DSN/PRICING_PATH/BATCH_SIZE 等、probe_ttl 派生补第三项、守卫补齐 |
src/polygateway/client.py |
M | from_settings 分支建 Redis 后端/PostgresRecorder/PricingTable;aclose 探测 async aclose |
pyproject.toml |
M | 核对既有 postgres extra(已存在于 pyproject.toml:22,勿重复添加);pytest addopts = -m "not slow"(见 T3);import-linter layers 契约已覆盖 backends 子包,无需改动 |
.env.example |
M | 回填全部新键 |
tests/contracts/conftest.py |
M | fixture 增 redis param + 哨兵时钟 |
tests/integration/test_redis_governance_time.py |
C | 时间语义 1:1 真实等待变体(slow) |
tests/integration/test_redis_cross_connection.py |
C | CHS 三个跨连接用例移植 + 联合验证 |
tests/integration/test_postgres_telemetry.py |
C | PG 幂等/降级/并发 |
tests/unit/test_redis_key_layout.py |
C | key 生成、秒毫秒换算、stats 读侧 clamp(纯函数部分) |
tests/unit/test_backpressure.py |
C | stall 四象限 + jitter + 记账降级 |
tests/unit/test_pricing.py |
C | 价格表/换算/未知模型/坏文件 |
tests/unit/test_embedding.py |
C | 分批/保序/归一化/维度校验/重试换源/遥测行 |
tests/unit/test_soak_corpus.py |
C | traces 还原器、不变量断言函数(纯函数) |
tests/e2e/test_embed_probe.py |
C | 真实网关 /embeddings 探测(e2e,允许记录降级) |
tools/soak/{corpus,scenarios,run_soak,scoreboard}.py |
C | harness 四组件(不被 src import) |
research-wiki/migrations/*.md |
M | embedding 迁移条目(VT adapter 说明) |
保真校验总则: T1/T2/T5 为迁移类任务,实现时逐段比对下列参考并在提交信息外的自查中确认语义一致;设计 §9 已声明的"替换/有意放弃/有意反转"之外,任何行为差异 = bug:
reference/CHSAnalyzer/app/coordination/scripts.py(全部 Lua)reference/CHSAnalyzer/app/coordination/limiter.py/provider_gate.pyreference/CHSAnalyzer/app/providers/governance.py:200-285(stall/jitter) T6-T9 非迁移(Postgres/pricing 零先例;Embedding 为按设计 §7.3 裁决表的重写,審計表即保真基准)。
T0 配置地基与依赖
文件: src/polygateway/config.py(M)、pyproject.toml(M)、.env.example(M)、tests/unit/test_config.py(M,增用例)
行为:
- 删除 config.py 对
PGW_LIMITER_BACKEND/PGW_BREAKER_BACKEND=redis的硬拒(现 254-255 行);取redis时校验REDIS_URL必在,否则 ValueError。 PGW_TELEMETRY_BACKEND白名单 →("sqlite","postgres","none");postgres时_require("PGW_TELEMETRY_PG_DSN");DSN 若含+asyncpg/+psycopg后缀在装配层剥为postgresql://。- 新可选键
PGW_PRICING_PATH(存GatewaySettings.pricing_path: Path | None)。 - probe_ttl 缺省派生式改为
max(2×最大源 timeout_s, cooldown_s, 最大源 timeout_s+5);新增装配守卫probe_ttl_s ≥ max(timeout_s)+5、核对既有守卫max(timeout_s) ≤ lease_ttl_s与stall_window_s ≥ max(ttft_timeout_s)(缺则补),违反 ValueError。 pyproject.toml: 仅核对——postgresextra 已存在(pyproject.toml:22),import-linter layers 契约自动涵盖backends.redis子包,两者均无需改动。- DSN 存
GatewaySettings.telemetry_pg_dsn: str | None(T6 消费此字段名)。 - 既有测试修订:
tests/unit/test_config.py:111-112显式PROBE_TTL_S=45(源 timeout=120)将触新守卫 → 改为合法值(≥125)并保留断言意图;另增守卫违反用例。.env实况已核:probe 走派生路径(240≥125),e2e 装配不受影响。
测试(先红): redis 后端键通过校验/缺 REDIS_URL 报错;postgres 缺 DSN 报错;DSN 后缀剥离;派生式三项取最大;守卫违反报错。
验证: conda run -n PolyGateway pytest tests/unit/test_config.py -v 全 PASS;make lint。
- 提交
feat: unlock redis/postgres backend config with assembly guards
T1 Redis 限流后端
文件: src/polygateway/backends/redis/{__init__,limiter.py}(C)、tests/unit/test_redis_key_layout.py(C)、tests/contracts/conftest.py(M,仅 limiter_factory 部分)
接口(跨任务消费,实际代码):
class RedisLimiter:
def __init__(self, *, scope: str, sources: dict[str, SourceConfig],
global_limits: GlobalLimits, redis: "redis.asyncio.Redis",
lease_ttl_s: float = 1500.0, poll_interval_s: float = 0.05) -> None: ...
@classmethod
def from_url(cls, url: str, **kwargs) -> "RedisLimiter": ...
# 契约方法与 InMemoryLimiter 完全同签名(ports.RateLimiter)
行为(逐段比对 CHS,设计 §2.2 保真点全数落实):
- Lua 常量:ACQUIRE(六闸同脚本,判据顺序与
>=/+est>语义照抄 scripts.py:6-34)、RELEASE(37-41)、SETTLE(44-48)、STATS(含 ZREMRANGEBYSCORE 清过期)、PROGRESS_MARK(60-67)、PROGRESS_AGE(70-78)。key 前缀pgw:limit:,布局pgw:limit:{GLOBAL|<scope>}:{<source>|}:lease|rpm:{win}|tpm:{win},窗口 id =redis.time() sec//60。 _RedisPermit:released/settled flag 幂等;settle 只动 TPM、INCRBY delta、落 acquire 窗口;source_stats读侧max(0, ·)clamp。- 限额 0 = 该闸不启用——这是对 CHS 的有意偏离,不是同款(核实:CHS scripts.py:18-23 六判据均为裸比较、无
>0守卫,0 在 CHS 语义下是全拒,CHS 靠传大数表示不限);M1 契约已冻结 0=禁用(_NO_GLOBAL = GlobalLimits(0,0,0)),故库版 Lua 显式新增limit > 0守卫。逐段比对时此处以设计 §9 勘误行为准,勿视为失真。 - RedisError:try_acquire/acquire/source_stats/progress_age_s →
GovernanceBackendError;settle/release →logger.warning降级(设计 §10)。mark_progress 也抛GovernanceBackendError(降级责任在中间件,T5)。 aclose()幂等关闭自建 client(from_url 路径);注入 client 不代关。
契约接入(红→绿证据,conftest 结构定案): 重构为单一 backend fixture(params=["memory","redis"],async fixture),clock 与两工厂都依赖它——memory 时 clock=FakeClock、工厂建内存版;redis 时 clock=哨兵 SkipClock(advance() 即 pytest.skip("redis 时间语义由 tests/integration/test_redis_governance_time.py 变体覆盖"))、工厂建 Redis 版(读 REDIS_URL,缺则 skip;每次构造唯一 scope=f"t{uuid4().hex[:8]}");redis param 的 fixture setup 先做分钟窗口 headroom 防抖(redis.time() 剩余窗口 <10s 则睡至翻滚,防 RPM/TPM 用例跨窗 flake,CHS _await_window_headroom 同款)。先提交 conftest 增参跑一次记录红(RedisLimiter 未实现 → 构造失败),实现后非时间用例全绿。
验证: conda run -n PolyGateway pytest tests/contracts/test_limiter_contract.py tests/unit/test_redis_key_layout.py -v → memory 全 PASS + redis 非时间用例 PASS + 时间用例 SKIP(2 个:test_lease_expiry_reclaims_slot、test_progress_marks_fresh)。
- 提交
feat: add redis six-gate rate limiter backend
T2 Redis 熔断后端
文件: src/polygateway/backends/redis/breaker.py(C)、tests/contracts/conftest.py(M,gate_factory 增 redis param)
接口: RedisGate(config: BreakerConfig, redis, scope: str) + from_url;契约方法同 ports.ProviderGate。
行为(逐段比对 provider_gate.py + scripts.py:80-245):
- 五条 Lua:TRY_ENTER(82-106)、RECORD_SUCCESS(110-143)、RECORD_FAILURE(148-199)、RELEASE_PROBE(203-225)、RETRY_AFTER(228-245)。HASH
pgw:gate:{scope}:{source},字段 state/epoch/failures/open_until/probe_until/probe_owner。 - 保真检查点:epoch 仅 record_failure 开断时 +1(scripts.py:180);探针/force_open failures 顶格(173-178);fencing 判据两分支(119-123/155-161);release_probe 置
open_until=now立即可接管(210-216);retry_after 全 clamp ≥0、空集合 Python 侧 ValueError(M1 契约);时间回拨 clamp。 GateDecision/GateUpdate由 Lua 返回数组组装,ms→s 换算后构造(冻结类型自带不变式校验即免费的移植正确性检查)。- RedisError:try_enter/retry_after_s →
GovernanceBackendError;record_success/record_failure/release_probe 抛GovernanceBackendError(降级在 T5 中间件)。
红→绿证据: gate_factory 增 redis param 先跑记录红;实现后非时间用例全绿(TestStateMachine 4 例 + TestEpochFencing.test_stale_epoch_write_rejected)。
验证: conda run -n PolyGateway pytest tests/contracts/ -v → redis param 下非时间用例 PASS、时间用例 SKIP 7 例(TestHalfOpenProbe 5 + test_stale_probe_owner_rejected + TestRetryAfter.test_retry_after_semantics;test_retry_after_takes_min_across_sources 无 advance,应 PASS)。
- 提交
feat: add redis circuit breaker gate with epoch fencing
T3 时间语义 1:1 真实等待变体(人类拍板:不缩放,真实量级)
文件: tests/integration/test_redis_governance_time.py(C)
注: 设计 §2.3 写"10 用例(限流 2+熔断 8)"为笔误——契约文件实际含 clock.advance 的用例是 9 个(熔断 7:test_stale_epoch_write_rejected 虽声明 clock 但无 advance,属非时间用例应 PASS),本表为准。slow 排除策略(全局): pyproject [tool.pytest.ini_options] 增 addopts = "-m \"not slow\""——pre-commit hook 与 make test 默认跳过 slow(否则每次提交等 15 分钟),显式 pytest -m slow 覆盖之(CLI 的 -m 后到覆盖 addopts);T12 显式跑全量 slow。
1:1 映射表(每行一个 @pytest.mark.slow 用例,断言与契约原用例同名行为;配置用真实量级:BreakerConfig(fail_threshold=3, cooldown_s=60.0, probe_ttl_s=120.0) 同契约 _CFG;limiter lease_ttl_s=30.0;等待 = asyncio.sleep(时长+1),唯 #2 例外精确睡 42s——其断言 41<age<43 有上界):
| # | 契约原用例(skip 于 redis param) | 变体断言要点 | 真实等待 |
|---|---|---|---|
| 1 | limiter test_lease_expiry_reclaims_slot |
泄漏 permit 30s 后槽位回收可再入 | ~31s |
| 2 | limiter test_progress_marks_fresh |
mark 后 age<5;等 42s 后 41<age<43 | ~42s |
| 3 | breaker test_cooldown_grants_single_probe |
冷却 60s 后唯一探针,第二进入者拒 | ~61s |
| 4 | breaker test_probe_success_closes |
探针成功 → CLOSED 可入 | ~61s |
| 5 | breaker test_probe_failure_reopens |
探针失败 → OPEN 拒入 | ~61s |
| 6 | breaker test_probe_lease_expiry_allows_takeover |
探针 120s 过期后接管 | ~182s |
| 7 | breaker test_release_probe_hands_back_and_idempotent |
归还幂等,下家立即拿探针 | ~61s |
| 8 | breaker test_stale_probe_owner_rejected |
死探针迟到写回被 fence | ~182s |
| 9 | breaker test_retry_after_semantics |
开路>0≤cooldown;到期 0;空集 ValueError | ~61s |
另加:RPM/TPM 分钟窗口翻滚用例(等待至下一分钟边界,移植 CHS tests/integration/test_redis_limiter.py:17-30 _await_window_headroom 防抖:剩余窗口 <10s 先等翻滚)。整文件预计 12-15 分钟。
完整性守卫: 文件头一个 meta 断言用例——解析两份契约测试源码(ast/正则)提取所有调用 clock.advance 的测试函数名集合,与本文件变体函数名集合(约定 test_variant_<原名>)比对相等;未来契约新增时间用例而变体漏配即此 meta 用例红(手维护清单达不到此效果)。
红→绿: 变体先写、RedisLimiter/RedisGate 已存在故直接跑;每例首跑必须真实观察等待(执行时禁止用 FakeClock 捷径)。
验证: conda run -n PolyGateway pytest tests/integration/test_redis_governance_time.py -v -m slow(全 PASS,时长 >600s 属预期)。
- 提交
test: add real-wait time-semantics variants for redis backends
T4 跨连接与多源联合验证
文件: tests/integration/test_redis_cross_connection.py(C)
行为(双连接池 = 多 worker 等价,人类已认可):
- 移植 CHS 用例(蓝本
reference/CHSAnalyzer/tests/integration/test_redis_limiter.py):2 个跨连接用例(:127 双连接池全局并发、:166 进度跨连接可见)+ 将test_global_rpm_cap(:105,CHS 原版为单连接)升级为跨连接版本(两池共享全局 RPM)。 - 联合验证:两个
GatewayClient(各自连接池、同 scope、文件内自建 ScriptedTransport——参照tests/unit/test_retry.pyFakeTransport 形态,库内无现成 MockTransport)+ 双源;断言 (a) 源 A 熔断开路后另一 client 的 try_enter 也被拒(状态共享);(b) 全局 RPM=N 时两 client 并发打满,遥测重算任意 60s 窗口 ≤N(不超配);(c) in-flight 取消 → lease 释放(source_stats.inflight归零)。 - Redis 掉线方向:对错误端口的 RedisLimiter/RedisGate,准入侧抛
GovernanceBackendError(fail-closed 集成证据)。
验证: conda run -n PolyGateway pytest tests/integration/test_redis_cross_connection.py -v 全 PASS。
- 提交
test: verify cross-connection shared governance state
T5 背压 stall + 记账侧降级(middleware/retry.py)
文件: src/polygateway/middleware/retry.py(M)、src/polygateway/middleware/ratelimit.py(M)、tests/unit/test_backpressure.py(C)
行为(比对 governance.py:200-285):
ratelimit.py的QuotaGate增透传方法async def progress_age_s(self) -> float(现仅有 try_acquire/stats/mark_progress,直接委托self._limiter.progress_age_s(),异常原样冒泡归准入侧语义)。- 循环入口记
entered_at = self._now()一次,不重置(CHS :207 口径)。 _on_no_runnablewait 分支:local_waited > stall_window_s and await self._quota.progress_age_s() > stall_window_s→raise AllSourcesExhausted(scope=self._scope, reason="stalled", retry_after_s=await self._breaker.retry_after_s(names), per_source_reasons=reasons)(构造函数 keyword-only,errors.py:72-80);否则await self._sleep(self._bp.poll_interval_s * (0.5 + 0.5 * self._rng()))(:283-285)。fail_fast 路径不变。- 记账侧降级:
record_success/record_failure/mark_progress/release_probe的调用点包try/except GovernanceBackendError: logger.warning(...),收敛为单个私有 helper_record_quietly(coro);不吞 CancelledError;_settle_and_release既有降级保持。 - 提取退避纯函数(T9 消费,本任务交付):实例方法
_backoff_delay(retry.py:254-259)提升为模块级def backoff_delay(policy: RetryPolicy, fails: int, exc: BaseException | None, rng: Callable[[], float]) -> float,RetryMW 内部改为调用它;行为零变化,既有 test_retry 回归即证据。
测试(先红,memory 后端 + FakeClock): stall 四象限(仅本地超窗不判死/仅全局超窗不判死/双超窗抛 stalled 且 reason 正确/未超窗继续 poll);jitter 界 [0.5p, p];记账降级——注入 record_success 抛 GovernanceBackendError 的 FakeGate,断言成功响应仍返回且有 warning;取消在 stall 等待中穿透。
验证: conda run -n PolyGateway pytest tests/unit/test_backpressure.py tests/unit/test_retry.py -v 全 PASS(既有 retry 测试回归)。
- 提交
feat: add dual-condition stall detection and quiet accounting degradation
T6 Postgres 遥测
文件: src/polygateway/telemetry/postgres.py(C)、src/polygateway/client.py(M:_build_telemetry 分支 + aclose 探测 async aclose)、tests/integration/test_postgres_telemetry.py(C)
接口:
class PostgresRecorder:
def __init__(self, dsn: str, *, pool: "asyncpg.Pool | None" = None) -> None: ...
async def record_llm_call(self, *, <18 字段同 ports.TelemetryRecorder>) -> None: ...
async def aclose(self) -> None: ...
行为: lazy 建池(首写);DDL CREATE TABLE IF NOT EXISTS llm_calls(18 列同 SQLite 语义,类型映射 TEXT/INT/DOUBLE PRECISION/BOOLEAN + created_at timestamptz DEFAULT now(),call_id TEXT PRIMARY KEY);写入 INSERT ... ON CONFLICT (call_id) DO NOTHING($n 占位);两级降级——建池/建表失败 warning 一次后池置 None 永久短路,单条 INSERT 失败逐条 warning 丢弃不降级;asyncpg 缺装 → 构造 ImportError 提示 pip install 'polygateway[postgres]'。
测试(先红,integration,读 .env PGW_TELEMETRY_PG_DSN,缺则 skip;表用 llm_calls_test_<uuid> 注入表名参数?——否:表名固定 llm_calls,测试前 DROP TABLE IF EXISTS,库本就 polygateway 专用): schema 18 列序;幂等(同 call_id 两写留首行);并发 50 写全落;坏 DSN 构造后写入不抛(结构性降级);单条超长参数导致的写失败不影响后续行(运行时降级)。
验证: conda run -n PolyGateway pytest tests/integration/test_postgres_telemetry.py -v 全 PASS。
- 提交
feat: add postgres telemetry recorder with two-tier degradation
T7 pricing
文件: src/polygateway/pricing.py(C)、src/polygateway/middleware/telemetry.py(M)、src/polygateway/client.py(M:from_settings 建表传入 Emitter)、tests/unit/test_pricing.py(C)
接口:
@dataclass(frozen=True)
class ModelPrice:
input_per_1m: float
output_per_1m: float
class PricingTable:
def __init__(self, prices: Mapping[str, ModelPrice]) -> None: ...
@classmethod
def from_file(cls, path: Path) -> "PricingTable": ...
def cost(self, model: str, prompt_tokens: int, completion_tokens: int) -> float | None: ...
行为: JSON 形态 {"<model>": {"input_per_1m": x, "output_per_1m": y}};文件缺失/JSON 坏/负单价 → from_file ValueError(装配 fail-loud);cost() 查不到 model → None + 该 model 仅首次 warning;TelemetryEmitter.__init__(recorder, *, pricing: PricingTable | None = None),_record 中 cost = 成功响应按 usage 换算、缓存命中 0.0、失败行 None;单一 helper 铁律不破(仍只此一个 record_llm_call 调用点,test_single_emitter_discipline 回归)。
测试(先红): 换算精度;未知 model None + 单次 warning;坏文件/负价拒绝;Emitter 成功行 cost 正确、cache_hit 行 0.0、失败行 None。
验证: conda run -n PolyGateway pytest tests/unit/test_pricing.py tests/unit/test_telemetry.py -v 全 PASS。
- 提交
feat: add pricing table and cost calculation in telemetry emitter
T8 Embedding 类型/端口/transport
文件: src/polygateway/types.py(M,只增)、src/polygateway/ports.py(M,只增)、src/polygateway/transports/openai_compat.py(M)、tests/unit/test_types.py(M,增用例)、tests/unit/test_embedding.py(C,transport 部分)
接口(设计 §7.2 冻结候选,照抄):
@dataclass(frozen=True)
class EmbeddingTransportResult:
vectors: list[list[float]]
dim: int
prompt_tokens: int
usage_source: str # measured | estimated
raw: dict[str, Any]
@dataclass(frozen=True)
class EmbeddingResponse:
vectors: list[list[float]]
dim: int
model: str
provider: str
prompt_tokens: int
usage_source: str
latency_ms: int
call_id: str
source_name: str
cost: float | None = None
@runtime_checkable
class EmbeddingTransport(Protocol):
async def embed(self, *, texts: list[str], source: SourceConfig,
call_id: str) -> EmbeddingTransportResult: ...
行为: OpenAICompatTransport.embed()——POST {base_url}/embeddings,body {"model": source.model, "input": texts};响应按 data[].index 排序取 embedding(GovDoc embedding.py:149-150 / VT :164-165 同款,保序硬约束);usage.prompt_tokens 有则 measured、无则 est_tokens estimated;HTTP 错误复用既有错误翻译表(401→SourceDead、429/5xx→Transient、400→RequestRejected);空 data/长度不匹配/向量维度互不一致 → ResultInvalidError。
测试(先红): frozen/只增字段;乱序 index 重排;usage 缺失 estimated;各错误码翻译;空响应 ResultInvalid。
验证: conda run -n PolyGateway pytest tests/unit/test_types.py tests/unit/test_embedding.py -v PASS。
- 提交
feat: add embedding types, port and openai-compat transport
T9 EmbeddingClient 治理循环 + 装配 + e2e 探测
文件: src/polygateway/embedding.py(C)、src/polygateway/config.py(M)、src/polygateway/__init__.py(M:导出)、tests/unit/test_embedding.py(M)、tests/e2e/test_embed_probe.py(C)
装配定案(不动 GatewaySettings 冻结形态): config.py 新增独立
@dataclass(frozen=True)
class EmbeddingSettings:
gateway: GatewaySettings # 复用 from_env(scope) 的源表/韧性/遥测/redis 字段
batch_size: int # {SCOPE}__BATCH_SIZE 必填,≥1
normalize: bool = False # {SCOPE}__NORMALIZE 可选(true/false)
expected_dim: int | None = None # {SCOPE}__EXPECTED_DIM 可选,≥1
@classmethod
def from_env(cls, scope: str = "EMBED", env=None) -> "EmbeddingSettings": ...
EMBED 专用键不进 GatewaySettings(LLM scope 不受影响);EmbeddingClient.from_settings(EmbeddingSettings, *, limiter=None, breaker=None, telemetry=None) + from_env 委托之(cache/structured 键对 embedding 无意义,装配时忽略)。
接口:
class EmbeddingClient:
def __init__(self, *, scope, sources, selector, limiter, breaker,
transport: EmbeddingTransport, retry: RetryPolicy,
backpressure: BackpressurePolicy, quota_full: str = "wait",
telemetry: TelemetryRecorder | None = None,
pricing: PricingTable | None = None,
batch_size: int, normalize: bool = False,
expected_dim: int | None = None,
now=time.monotonic, sleep=asyncio.sleep, rng=random.random) -> None: ...
async def embed(self, texts: list[str], *, session_id=None,
parent_call_id=None) -> EmbeddingResponse: ...
@classmethod
def from_env(cls, scope: str = "EMBED", **注入同 GatewayClient) -> "EmbeddingClient": ...
async def aclose(self) -> None: ...
行为: 空输入即返空响应零请求;按 batch_size 切批,批间串行;每批 = 完整治理调用:选源(selector+cooldown memo)→ try_enter → try_acquire → transport.embed → 成功 record_success+settle+mark_progress / 失败按四分类 record_failure(SourceDead force_open)换源退避——退避公式复用 retry.py 提取的模块级纯函数 backoff_delay(policy, fails, exc, rng)(T5 中提取,此处 import);stall 双条件与 quota_full 语义同 chat;每批遥测一行(messages=texts 截断摘要 JSON、response="<vectors n=N dim=D>"、thinking=""、completion_tokens=0),cost 经 pricing;聚合响应 vectors 拼接、prompt_tokens 求和、latency_ms 求和、usage_source 任一 estimated 则 estimated、call_id 取首批;normalize=True 时纯 Python L2 归一化(max(norm,1e-12) 防除零,VT embedding.py:167-170 语义);expected_dim 不符 → ResultInvalidError。取消穿透:批间/等待/请求中取消即上抛,in-flight permit/probe 在 finally 释放。
测试(先红,ScriptedEmbeddingTransport): 分批数与保序;归一化数值;expected_dim 违约;TransientError 重试换源、SourceDead 熔断换源;批失败时已完成批不重发(整调用失败上抛,已耗批次遥测可见);取消释放;遥测行字段与 usage_source 聚合。
e2e 探测(人类默认认可): test_embed_probe.py 对 .env MiniMax 源发一次真实 /embeddings 请求(无 EMBED 源配置则用 LLM 源 base_url);200 → 断言向量返回并落 tests/outputs/embedding/;404/不支持 → pytest.skip 并把响应记录进 outputs(记录降级证据)。
验证: conda run -n PolyGateway pytest tests/unit/test_embedding.py -v PASS;e2e 单独跑记录结果。
- 提交
feat: add governed embedding client with batching
T10 harness:语料与场景
文件: tools/soak/__init__.py、tools/soak/corpus.py、tools/soak/scenarios.py(C)、tests/unit/test_soak_corpus.py(C)
行为(findings §1/§2 为准):
corpus.py:load_trace_chains(harness_db) -> list[list[Messages]](P1:traces 表还原 system+多轮累积链,session/parent 链路字段保留);load_replay_payloads(telemetry_db) -> list[Messages](P2:376 条 messages 原样);assemble_frame_request(frames_dir, n_frames)(1-4/5-6 帧两档 base64 组装);chs_image_request(images_dir, structured_schema)(单图+12 字段固定 instruction,pydantic schema 在 tools 内定义——业务 schema 不入 src);全部纯函数可 unit 测。scenarios.py: P1-P6 各一个 async 生成器,yield(kind, kwargs)(kind=chat/embed);P4 = P3 流量 20-30% 重复 messages + salt 对照组;P5 = P3 打到含故障源的多源池(故障源配置 = findings §3 四类,坏 key/黑洞 IP/紧看门狗 1s/0.5s/紧闸 rpm=5 由 .env EMBED 之外的独立 scopeSOAK__*提供);P6 按签字比例 P1 10%/P2 20%/P3 50%(含 3/10 缓存双向)/P5 20% 加权混合。 测试(先红): trace 还原链长与角色序;帧组装档位;比例加权抽样统计近似(容差 ±5%)。 验证:conda run -n PolyGateway pytest tests/unit/test_soak_corpus.py -vPASS。- 提交
feat: add soak corpus loaders and scenario generators
T11 harness:入口与记分板
文件: tools/soak/run_soak.py、tools/soak/scoreboard.py(C)、tests/unit/test_soak_corpus.py(M,增不变量函数用例)
行为:
run_soak.py: argparse--scenario P1..P6 --budget-calls N --budget-tokens M --workers K --concurrency C --run-id(缺省 run-id=时间戳);启动守卫——REDIS_URL必须以/3结尾否则拒跑、PGW_TELEMETRY_PG_DSN(若配)路径库名必须polygateway否则拒跑;开跑 FLUSHDB(仅 db3)+cache_namespace=run_id+ 回放请求cache_salt=run_id;预算硬顶(签字值:全程 token ≤2 亿,单场景 calls 上限见设计 §8.1)任一命中优雅停(等 in-flight 收尾);--workers>1用multiprocessing起真实子进程,预算按 worker 均分,各 worker 独立遥测 db 文件data/soak/telemetry_<run_id>_<w>.db;全局并发/RPM 保护(签字值 100/600)配在 SOAK scope 源上。scoreboard.py: 断言 findings §4 六条硬不变量,产出报告tests/outputs/soak/<run_id>.md(结构化 Markdown)。输入口径分三类:不变量 2/3/5/6 = 纯函数(输入合并后的遥测行迭代器,unit 可测);不变量 1(记账归零)= 跑后经活 Redis 连接查source_stats.inflight==0且 gate 可再准入(集成性检查,unit 用注入的 fake limiter/gate 测判定逻辑);不变量 4(RSS 平稳)= 纯函数(输入 run_soak 周期采样落盘的 RSS json)。 测试(先红): 六条不变量函数各造一组违例数据断言报错、一组合规数据通过;db3 守卫拒非 /3 连接串。 验证:conda run -n PolyGateway pytest tests/unit/test_soak_corpus.py -vPASS;conda run -n PolyGateway python tools/soak/run_soak.py --scenario P3 --budget-calls 6 --budget-tokens 50000 --workers 2真实小跑通(冒烟,必须带 --workers 2——多进程验收口径的执行证据,人类认可 §13.5 的第二条腿;记入 outputs)。- 提交
feat: add soak runner with budget guards and invariant scoreboard
T12 收尾与独立验证
文件: .env.example(M)、research-wiki/migrations/video-tree-trm5.md 与 govdoc-saas.md(M:embedding 迁移条目——VT 需同步壳+ndarray adapter,GovDoc 替换 OpenAICompatEmbedding 映射)、research-wiki/ROADMAP.md(M:M2 状态)
行为: .env.example 回填全部新键(PGW_TELEMETRY_PG_DSN、PGW_PRICING_PATH、EMBED__、SOAK__ 示例,含"严禁指向在用库"注释;REDIS_URL 示例改为 /3 并注明 db3 纪律——现示例是 /0);ARCHITECTURE §7.3 记账侧勘误已随设计批准提交(commit b165c2a),本任务无需重做,verifier 勿误报缺失;make ci 全绿 + 显式跑全量 slow(addopts 默认排除);覆盖率 ≥80%;派全新上下文 verifier subagent(只读)对照设计逐条验收(verification-before-completion),Critical/Important 清零后方可声称完成;按 findings §2 预算跑一轮 P3 中等规模(≤500 calls,--workers 2)真实验收记入 outputs。
验证: conda run -n PolyGateway make ci 输出 0 failed;conda run -n PolyGateway pytest -m slow -v 全 PASS(≥12 分钟属预期);verifier 报告落 findings(若有问题)。
- 提交
docs: backfill env template and migration notes for M2
依赖与并行
T0 → T1 → T2 → T3/T4(可并行)→ T5;T6/T7 依赖 T0 可与 T1-T5 并行;T8 → T9(依赖 T5 的 backoff 提取);T10 → T11(依赖 T1-T9 全部);T12 最后。全程在 feature/m2-distributed 分支,每任务一提交。
计划自查清单(写完即核)
- 设计 §1-§8 每节均有对应任务(§1→T0、§2→T1/T3、§3→T2/T3、§4→T5、§5→T6、§6→T7、§7→T8/T9、§8→T10/T11);签字值(token 2 亿/并发 100/RPM 600/P6 比例)已写死进 T11;三口径(时间变体/多进程/记账降级)分别落 T3/T11/T5。
- 无 placeholder;跨任务接口(RedisLimiter/PostgresRecorder/PricingTable/Embedding 三类型/backoff_delay 提取)均有实际签名。
- 保真校验:T1/T2/T5 标注 CHS 文件行号并要求逐段比对;T8/T9 以设计 §7.3 裁决表为基准。