Files
PolyGateway/research-wiki/plans/2026-07-20-m2-distributed-plan.md

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 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_sstall_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 部分)

接口(跨任务消费,实际代码):

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_slottest_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 等价,人类已认可):

  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.pyQuotaGate 增透传方法 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_sraise 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)

接口:

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__.pytools/soak/corpus.pytools/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.pytools/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>1multiprocessing 起真实子进程,预算按 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 -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.mdgovdoc-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 裁决表为基准。