# 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 env `PolyGateway`;真实 Redis = `.env` `REDIS_URL`(**db3 专用**);真实 Postgres = `.env` `PGW_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.py` - `reference/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,增用例) **行为**: 1. 删除 config.py 对 `PGW_LIMITER_BACKEND/PGW_BREAKER_BACKEND=redis` 的硬拒(现 254-255 行);取 `redis` 时校验 `REDIS_URL` 必在,否则 ValueError。 2. `PGW_TELEMETRY_BACKEND` 白名单 → `("sqlite","postgres","none")`;`postgres` 时 `_require("PGW_TELEMETRY_PG_DSN")`;DSN 若含 `+asyncpg`/`+psycopg` 后缀在装配层剥为 `postgresql://`。 3. 新可选键 `PGW_PRICING_PATH`(存 `GatewaySettings.pricing_path: Path | None`)。 4. 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。 5. `pyproject.toml`: 仅核对——`postgres` extra 已存在(pyproject.toml:22),import-linter layers 契约自动涵盖 `backends.redis` 子包,两者**均无需改动**。 6. DSN 存 `GatewaySettings.telemetry_pg_dsn: str | None`(T6 消费此字段名)。 7. **既有测试修订**: `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 部分) **接口**(跨任务消费,实际代码): ```python 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|}:{|}: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**——其断言 `410≤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 等价,人类已认可): 1. 移植 CHS 用例(蓝本 `reference/CHSAnalyzer/tests/integration/test_redis_limiter.py`):**2 个跨连接用例**(:127 双连接池全局并发、:166 进度跨连接可见)+ 将 `test_global_rpm_cap`(:105,CHS 原版为单连接)**升级为跨连接**版本(两池共享全局 RPM)。 2. 联合验证:两个 `GatewayClient`(各自连接池、同 scope、文件内自建 ScriptedTransport——参照 `tests/unit/test_retry.py` FakeTransport 形态,库内无现成 MockTransport)+ 双源;断言 (a) 源 A 熔断开路后另一 client 的 try_enter 也被拒(状态共享);(b) 全局 RPM=N 时两 client 并发打满,遥测重算任意 60s 窗口 ≤N(不超配);(c) in-flight 取消 → lease 释放(`source_stats.inflight` 归零)。 3. 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): 1. `ratelimit.py` 的 `QuotaGate` 增透传方法 `async def progress_age_s(self) -> float`(现仅有 try_acquire/stats/mark_progress,直接委托 `self._limiter.progress_age_s()`,异常原样冒泡归准入侧语义)。 2. 循环入口记 `entered_at = self._now()` 一次,**不重置**(CHS :207 口径)。 3. `_on_no_runnable` wait 分支:`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 路径不变。 4. 记账侧降级:`record_success`/`record_failure`/`mark_progress`/`release_probe` 的调用点包 `try/except GovernanceBackendError: logger.warning(...)`,收敛为单个私有 helper `_record_quietly(coro)`;不吞 CancelledError;`_settle_and_release` 既有降级保持。 5. **提取退避纯函数**(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) **接口**: ```python 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_` 注入表名参数?——否:表名固定 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) **接口**: ```python @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 形态 `{"": {"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 冻结候选,照抄): ```python @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 新增独立 ```python @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 无意义,装配时忽略)。 **接口**: ```python 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=`""`、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 之外的独立 scope `SOAK__*` 提供);P6 按签字比例 P1 10%/P2 20%/P3 50%(含 3/10 缓存双向)/P5 20% 加权混合。 **测试(先红)**: trace 还原链长与角色序;帧组装档位;比例加权抽样统计近似(容差 ±5%)。 **验证**: `conda run -n PolyGateway pytest tests/unit/test_soak_corpus.py -v` PASS。 - [ ] 提交 `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__.db`;全局并发/RPM 保护(签字值 100/600)配在 SOAK scope 源上。 - `scoreboard.py`: 断言 findings §4 六条硬不变量,产出报告 `tests/outputs/soak/.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 -v` PASS;`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 裁决表为基准。