docs: add M2 distributed implementation plan
This commit is contained in:
@@ -0,0 +1,304 @@
|
||||
# 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 + 记账侧降级捕获 |
|
||||
| `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 | extra `postgres = ["asyncpg>=0.29"]`;import-linter 契约补 redis 后端模块 |
|
||||
| `.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`: extra `postgres=["asyncpg>=0.29"]` 并入 `all`;import-linter 独立实现层清单加 `polygateway.backends.redis`。
|
||||
|
||||
**测试(先红)**: 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|<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 = 该闸不启用(与内存版一致,Lua ARGV 传 0 时跳过该闸——CHS 同款,比对 scripts.py 闸判据的 `>0` 守卫)。
|
||||
- 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 `limiter_factory` params 增 `"redis"`:读 `REDIS_URL`(缺则 `pytest.skip`),每次构造用唯一 `scope=f"t{uuid4().hex[:8]}"`;`clock` fixture 对 redis param 返回哨兵 `SkipClock`(`advance()` 即 `pytest.skip("redis 时间语义由 tests/integration/test_redis_governance_time.py 变体覆盖")`)。**先提交 conftest 增参跑一次记录红**(RedisLimiter 未实现 → collection/构造失败),实现后非时间用例全绿。
|
||||
|
||||
**验证**: `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)
|
||||
|
||||
**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)`):
|
||||
|
||||
| # | 契约原用例(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 断言用例——收集 redis param 下契约 skip 的用例名集合(硬编码上表 9 名),与本文件变体函数名集合比对,防未来契约新增时间用例而变体漏配。
|
||||
**红→绿**: 变体先写、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:105-188`):双连接池共享全局并发 / 共享全局 RPM / 进度跨连接可见。
|
||||
2. 联合验证:两个 `GatewayClient`(各自连接池、同 scope、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)、`tests/unit/test_backpressure.py`(C)
|
||||
|
||||
**行为**(比对 governance.py:200-285):
|
||||
1. 循环入口记 `entered_at = self._now()` 一次,**不重置**(CHS :207 口径)。
|
||||
2. `_on_no_runnable` wait 分支:`local_waited > stall_window_s and await self._quota.progress_age_s() > stall_window_s` → `raise AllSourcesExhausted(scope, reason="stalled", retry_after_s=await self._breaker.retry_after_s(names), per_source_reasons=reasons)`;否则 `await self._sleep(self._bp.poll_interval_s * (0.5 + 0.5 * self._rng()))`(:283-285)。fail_fast 路径不变。
|
||||
3. 记账侧降级:`record_success`/`record_failure`/`mark_progress`/`release_probe` 的调用点包 `try/except GovernanceBackendError: logger.warning(...)`,收敛为单个私有 helper `_record_quietly(coro)`;不吞 CancelledError;`_settle_and_release` 既有降级保持。
|
||||
|
||||
**测试(先红,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_<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)
|
||||
|
||||
**接口**:
|
||||
```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 形态 `{"<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 冻结候选,照抄):
|
||||
```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:`{SCOPE}__BATCH_SIZE` 必填、`{SCOPE}__NORMALIZE`/`{SCOPE}__EXPECTED_DIM` 可选)、`src/polygateway/__init__.py`(M:导出)、`tests/unit/test_embedding.py`(M)、`tests/e2e/test_embed_probe.py`(C)
|
||||
|
||||
**接口**:
|
||||
```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=`"<vectors n=N dim=D>"`、thinking=""、completion_tokens=0),cost 经 pricing;聚合响应 vectors 拼接、tokens 求和、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_<run_id>_<w>.db`;全局并发/RPM 保护(签字值 100/600)配在 SOAK scope 源上。
|
||||
- `scoreboard.py`: 读全部 worker 遥测库合并,断言 findings §4 六条硬不变量(1 记账归零:结束后 source_stats inflight==0 且 gate 可再准入;2 行数==请求数±取消、call_id 唯一;3 按时间戳重算任意 60s 窗口单源请求 ≤RPM;4 RSS 首末差 <阈值(run_soak 周期采样自身 RSS 落 json);5 P3 structured 最终成功率 ≥基线(首跑建立基线文件);6 成本/延迟分位),产出报告 `tests/outputs/soak/<run_id>.md`(结构化 Markdown)。不变量断言实现为纯函数(输入行迭代器),unit 可测。
|
||||
**测试(先红)**: 六条不变量函数各造一组违例数据断言报错、一组合规数据通过;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 5 --budget-tokens 50000` 真实小跑通(冒烟,记入 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__* 示例,含"严禁指向在用库"注释);`make ci` 全绿(含 slow 集成;运行环境需 REDIS_URL 与 PG DSN);覆盖率 ≥80%;派**全新上下文 verifier subagent**(只读)对照设计逐条验收(verification-before-completion),Critical/Important 清零后方可声称完成;按 findings §2 预算跑一轮 P3 中等规模(≤500 calls)真实验收记入 outputs。
|
||||
**验证**: `conda run -n PolyGateway make ci` 输出 0 failed;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 裁决表为基准。
|
||||
Reference in New Issue
Block a user