22 KiB
M2 分布式里程碑设计:Redis 治理后端 + 背压 + Postgres 遥测 + pricing + Embedding + 压测 harness
状态: 草案(待 Claude 自审 → 独立 subagent 审 → 人类门) 依据: ARCHITECTURE.md §7.3/§7.4/§7.8、ROADMAP §3(Q3 已拍板纳入 Embedding)、findings/2026-07-20-m2-soak-workload.md、M1 冻结契约(designs/2026-07-20-m1-core-design.md)、三份 reference 逐字调研(CHS 协调层 / Embedding 两版 / Postgres+pricing) 硬约束: M1 冻结的公共签名与双后端契约(
ports.py的RateLimiter/Permit/ProviderGate/GateDecision/GateUpdate/TelemetryRecorder18 字段)不改动;Redis 后端必须通过tests/contracts/同一套契约测试。
0. 范围与非目标
范围(七项 + harness): ① Redis 限流六道闸 ② Redis 熔断(epoch fencing)③ 多源 × Redis 联合验证 ④ 背压 stall 判定 ⑤ telemetry/postgres.py ⑥ pricing.py 成本入遥测 ⑦ Embedding 客户端 ⑧ 真实数据压测 harness(tools/,不入 pytest 门)。
非目标: OCR(M3)、音频、SDK transport、泛化中间件洋葱使其类型无关(见 §7 方案 G1 否决)、embedding 结果缓存(下游索引本就持久化向量,YAGNI)。
新增依赖(均须人类确认,随本设计一并批): redis(已是 optional extra,M1 缓存已用)、asyncpg → 新 extra postgres。Embedding 走 httpx,零新依赖。
1. 装配变化(config + 工厂)
| 键 | 变化 |
|---|---|
PGW_LIMITER_BACKEND / PGW_BREAKER_BACKEND |
解禁 redis(删除 config.py:254-255 硬拒);取 redis 时强校验 REDIS_URL 存在 |
PGW_TELEMETRY_BACKEND |
白名单扩为 sqlite/postgres/none;postgres 时强校验新键 PGW_TELEMETRY_PG_DSN |
PGW_PRICING_PATH |
新增,可选;指向价格表 JSON 文件;缺省 = 无 pricing,cost 恒 None(现状) |
EMBED__{PROVIDER}__{N}__{FIELD} |
Embedding 源沿用既有多源命名约定,EmbeddingClient.from_env(scope="EMBED") 装配 |
{SCOPE}__BATCH_SIZE |
Embedding scope 专用,必填(embedding 分批是行为关键,不设默认) |
装配期守卫(ARCHITECTURE §7.3/§7.4 契约补强,后端无关,对 memory/redis 一体执法):max(timeout_s) ≤ lease_ttl_s、probe_ttl_s ≥ max(timeout_s) + 5、stall_window_s ≥ max(ttft_timeout_s)(未配 ttft 的源跳过)。违反 → 装配报错拒绝启动。M1 已实现的部分保持,缺的补齐(plan 核对)。
共享语义不变:多逻辑角色共享全局闸 = 显式把同一 RedisLimiter/RedisGate 实例传入多次 from_env(limiter=..., breaker=...);Redis 后端天然跨进程共享(key 按 scope+source),同 scope 的多进程 worker 无需显式传实例即共享状态——这是与内存版唯一的行为差异,属 Redis 后端的目的本身。
2. Redis 限流六道闸(backends/redis/limiter.py)
2.1 方案对比
| 方案 | 内容 | 裁决 |
|---|---|---|
| L1 单条 Lua 逐字移植 | CHS scripts.py ACQUIRE/RELEASE/SETTLE/STATS/PROGRESS_* 六脚本原样移植,读-判-占一次 EVAL 原子完成 |
推荐:语义已被 CHS 生产验证 + 契约测试钉死;保真校验成本最低 |
| L2 WATCH/MULTI 事务 | redis-py 乐观锁重写 | 否决:多往返、竞态重试风暴、六闸联合判定难以原子化 |
| L3 每闸独立脚本 + 补偿 | 六个小脚本,失败回滚已占闸 | 否决:原子性碎裂,部分占用的补偿路径本身可能失败 |
2.2 关键语义(逐字保真点,plan 中逐条比对 CHS 行号)
- Key 布局(前缀 CHS
cclimit:→pgw:limit:):并发 = ZSETpgw:limit:{GLOBAL|scope}:{...}:lease(member=lease_id, score=过期时刻 ms);RPM/TPM = string 计数...:rpm:{win}/...:tpm:{win},窗口后缀由 Python 侧预生成。 - 窗口 id 用 Redis 服务器时钟:
redis.time() → int(sec)//60,多进程口径统一;Lua 内部TIME的 ms 只用于 lease 过期判定。时钟三源分工不得合并(窗口=服务器秒、lease=服务器 ms、stall 本地等待=monotonic)。 - 检查顺序与判据:全局并发→单源并发→全局 RPM→单源 RPM→全局 TPM→单源 TPM;并发/RPM 用
>=(占后即满),TPM 用used + est >(预扣后是否超)。拒绝零副作用(仅顺带ZREMRANGEBYSCORE清过期 lease)。 - 租约双保险:lease score 过期惰性清理 + 整 ZSET
PEXPIRE防僵尸 key;RPM/TPM keyEXPIRE 120s兜底。 - settle 只结算 TPM,
INCRBY delta(可为负,与内存版max(0,·)的钳位差异由契约测试的<=容忍断言吸收),落 acquire 时刻的窗口而非当前窗口;RPM/并发不退款——熔断期白烧 RPM 由 M1 已交付的源冷却备忘缓解(retry.py:168-170)。 - release/settle 幂等:Python 侧 flag + Lua ZREM 天然幂等;settle 仅成功后置位,失败可重试。
- Redis 异常 →
GovernanceBackendError(准入侧 fail-closed 报错不放行);settle/release 释放侧失败降级 warning 且不掩盖主异常(M1 勘误既定方向,ports 契约不变)。 - 进度键:
pgw:limit:GLOBAL:{scope}:progress_ms,SET 服务器 ms +EXPIRE 3600(TTL 远大于 stall_window,防进度键过期被误判"从未进展");progress_age_s:键缺失 →inf,时钟回拨 clamp 0。 - 秒(契约)↔ 毫秒(Redis 内部)换算是后端私事,不进端口。
2.3 契约测试接入
| 方案 | 内容 | 裁决 |
|---|---|---|
| T1 fixture 增参 + 真实 Redis | tests/contracts/conftest.py 的 limiter_factory/gate_factory params 增 "redis";无 REDIS_URL 时该 param skip;每 test 唯一 scope(uuid)隔离 |
推荐:M1 预留的接入方式,测试体零改动 |
| T2 fakeredis | 进程内模拟 | 否决:不执行真实 Lua,违反"Redis 测试用真实 Redis"规约 |
时间语义用例(租约过期、窗口翻滚、进度 age 推进)FakeClock 对 Redis 不可用 → 该三类用例在 redis param 下改小 TTL + 真实等待变体(独立 integration 用例,移植 CHS _await_window_headroom 分钟翻滚防抖),契约文件本体不改。测试固定用 db3。
3. Redis 熔断(backends/redis/breaker.py)
方案同构于 §2(L1 逐字移植 CHS provider_gate.py 五条 Lua,否决理由同,不重复列)。Key:每源一个 HASH pgw:gate:{scope}:{source},字段 state/epoch/failures/open_until/probe_until/probe_owner。
保真点(与 M1 内存版契约测试逐条对应,Redis 版必须同绿):
- epoch 仅在进入 OPEN 时 +1(record_failure 开断路径);try_enter 发探针、record_success、release_probe 均不递增。
- fencing 判据:普通写回
state==closed && epoch==entry.epoch;探针写回state==half_open && epoch匹配 && owner==entry.probe_owner;不匹配 →applied=False快照返回。allowed=False的决定伪造写回在 Python 侧拒绝(ValueError,与内存版一致)。 - 探针租约 =
probe_until绝对 ms + 惰性重发(无看门狗);release_probe置open_until=now使下家立即接管。 - 探针失败/
force_open→ failures 顶格 threshold 后开断;普通失败累加。 retry_after_s(sources)取集合最小等待,clamp ≥0;空集合 ValueError(M1 契约)。BreakerConfig.probe_ttl_s保持显式配置(M1 冻结形态),CHS 的"container 自动算max(timeout)+5"改为 §1 装配守卫——方案对比:自动计算(改冻结配置形态,否决)vs 显式 + 守卫(推荐,校验同效且不动公共类型)。
4. 背压 stall(middleware/retry.py 增量)
| 方案 | 内容 | 裁决 |
|---|---|---|
S1 在 _on_no_runnable wait 分支内实现 |
首次无可运行源时记 entered_at(本地 monotonic);双条件 local_waited > stall_window_s && progress_age_s() > stall_window_s 同时成立 → 抛 AllSourcesExhausted(reason="stalled");否则 sleep poll_interval_s * (0.5+0.5*rng)(jitter 防惊群,CHS governance.py:283-285) |
推荐:挂接点已在 M1 预留,不新增洋葱层 |
| S2 独立 BackpressureMW | 新中间件层 | 否决:改冻结层序;stall 判定与选源循环共享状态,拆层反而耦合 |
细则:entered_at 在每次 chat() 调用的重试循环内首次进入等待时记录,拿到 permit 即重置(等待是"连续无进展"语义,非累计);双条件缺一不判死(本地 monotonic 与 Redis 服务器时钟刻意不混用,CHS governance.py:270-281);reason="stalled" 已在 M1 错误模型 5 值枚举内,retry_after_s 取 breaker.retry_after_s(全部源)。mark_progress 调用点 M1 已就位(retry.py:213 成功即标)。行为对 memory/redis 两后端一致(progress_age_s 是端口方法)。
5. Postgres 遥测(telemetry/postgres.py,extra postgres)
参考现实:三项目均无 Postgres 遥测先例;工程蓝本取 GovDoc PostgresTaskStore(asyncpg、$n 占位、CREATE TABLE IF NOT EXISTS、ON CONFLICT DO NOTHING),但其"失败冒泡"方向与遥测铁律相反,降级方向反向处理。
| 方案 | 内容 | 裁决 |
|---|---|---|
| P1 DSN 自建池(lazy) | PostgresRecorder(dsn),首次写入时 asyncpg.create_pool;建池/建表/写入任何失败 → logger.warning 一次性降级(置池 None 短路后续写),不冒泡 |
推荐:与 SQLiteRecorder 对称,from_settings 只需 DSN 字符串 |
| P2 注入 asyncpg.Pool | 池由业务创建传入 | 保留为构造函数可选参数(pool= 优先于 dsn 自建),满足"构造函数全量注入"路线;不作为 from_env 路径 |
| P3 复用 SQLAlchemy | — | 否决:重依赖,违依赖极简 |
细则:schema = 18 列同名 + created_at timestamptz DEFAULT now(),call_id TEXT PRIMARY KEY;幂等 ON CONFLICT (call_id) DO NOTHING;asyncpg 原生异步,无 to_thread/Lock;DSN 若带 SQLAlchemy 风格 +asyncpg 后缀先剥离(GovDoc verify 脚本教训);关闭走 aclose()(async),GatewayClient.aclose 增加对 recorder async close 的探测(内部行为,非冻结签名)。降级语义:构造不连库(lazy),Postgres 不可用时业务调用零感知(warning 一次,后续静默),与"缓存/遥测静默降级"铁律一致。写失败不重试(遥测丢一条 < 拖垮调用)。
6. pricing(pricing.py)
参考现实:三项目零先例,从头设计。
| 方案 | 内容 | 裁决 |
|---|---|---|
| B1 JSON 价格表文件 + 注入 | PricingTable.from_file(path)(PGW_PRICING_PATH)或 PricingTable(dict) 直接注入;条目 model → {input_per_1m, output_per_1m} |
推荐:价格随网关计费变化,由使用方维护 |
| B2 库内置默认价格表 | 硬编码常见模型单价 | 否决:必然过时的默认值掩盖真实成本(P5);实验室走中转网关,计费非官方价 |
| B3 env 平铺单价键 | PRICING__<model>__INPUT |
否决:model 名含 ./-,env 键名不友好 |
细则:cost = prompt/1e6*input + completion/1e6*output(float,币种由使用方全表统一口径,库不设币种字段——18 字段冻结);查不到 model → cost=None + 每 model 仅首次 warning(防日志风暴),不阻塞调用;计算点 = TelemetryEmitter(cost=None 的唯一现存占位点 middleware/telemetry.py:146),Emitter 构造增可选 pricing: PricingTable | None——遥测单一 helper 铁律不破。缓存命中行 cost=0(未产生新调用)。文件解析失败 → 装配报错(配置类失败 fail-loud,非运行时降级)。
7. Embedding 客户端(embedding.py + transports/openai_compat.py 增量)
7.1 方案对比(治理栈复用方式)
| 方案 | 内容 | 裁决 |
|---|---|---|
| G1 泛化中间件洋葱 | RetryMW/遥测/洋葱全部泛型化为请求类型无关 | 否决:动 M1 冻结核心,收益仅一个新调用形态,典型 gold-plating |
G2 独立 EmbeddingClient + 复用后端与算法件 |
新类持自己的精简治理循环(选源→冷却备忘→熔断门→限流 permit→transport→错误分类→退避),直接复用:RateLimiter/ProviderGate 端口及两种后端、errors.py 四分类、退避公式、SourceCooldownMemo、TelemetryEmitter、SourceConfig/选源器 |
推荐:零改冻结面;循环逻辑与 RetryMW 存在有限重复,以"共享算法件、循环骨架各自持有"为界(chat 循环含流式/结构化/缓存分支,embedding 循环无,强行合一才是复制) |
| G3 塞进 chat 洋葱 | embedding 请求伪装 ChatRequest | 否决:messages/stream/structured 全不适用,类型欺骗 |
7.2 新公共签名(冻结候选,过人类门后与 M1 同等约束)
@dataclass(frozen=True)
class EmbeddingResponse:
vectors: list[list[float]] # 与输入等长、保序
dim: int
model: str
provider: str
prompt_tokens: int
usage_source: str # measured | estimated
latency_ms: int
call_id: str
source_name: str
cost: float | None = None
class EmbeddingTransport(Protocol): # ports.py 新增
async def embed(self, *, texts: list[str], source: SourceConfig,
call_id: str) -> EmbeddingTransportResult: ...
# EmbeddingTransportResult(types.py 新增, frozen): vectors/dim/prompt_tokens/usage_source/raw
class EmbeddingClient:
async def embed(self, texts: list[str], *, session_id: str | None = None,
parent_call_id: str | None = None) -> EmbeddingResponse: ...
# from_env(scope="EMBED", limiter=..., breaker=..., telemetry=...) 与 GatewayClient 工厂对称
7.3 行为裁决(两版审计的分歧点)
| 分歧 | GovDoc | VT | 库裁决 |
|---|---|---|---|
| 同步/异步 | httpx async | 同步 SDK | async + httpx(纯 asyncio 中立铁律);VT 迁移侧自包同步壳 |
| 返回类型 | list[list[float]] | 归一化 np.ndarray | list[list[float]](核心不依赖 numpy);ndarray/tensor 转换留 VT 业务侧 adapter(记入迁移文档) |
| L2 归一化 | 不做 | 强制做 | 构造参数 normalize: bool(默认 False;VT 装配传 True;纯 Python 实现,防除零 max(norm,1e-12) 保 VT 语义) |
| 分批 | 内建 batch_size 切片 | 不分批整发 | 内建必填 batch_size;批间串行(并发交 gather_bounded);每批 = 一次完整治理调用(独立 permit/熔断记账/遥测行) |
| 重试 | 自研退避(无上限无 jitter) | 零重试 | 有意放弃两者,统一走库退避公式(base·2^n 封顶 + jitter,§7.2)与四分类驱动换源 |
| 维度校验 | 校验不截断 | 不校验 | 可选 `expected_dim: int |
| 输入形态 | 仅 list | str 或 list | 仅 list[str](显式优于隐式);空 list 返回空响应不发请求 |
| usage | on_usage 回调 | 丢弃 | 有意放弃回调,usage 直接入遥测(prompt_tokens=measured/estimated,completion_tokens=0),cost 经 pricing 换算 |
| index 排序 | 做 | 做 | 保留(响应按 index 重排保序) |
多批聚合:EmbeddingResponse 为全批合并(vectors 拼接、prompt_tokens 求和、latency 求和、call_id 取首批);遥测按批逐行(每批独立 call_id,parent_call_id 透传调用方值)。协议:POST {base_url}/embeddings(openai 兼容,与 chat 同 transport 类,错误翻译复用 §6.2 表——401→SourceDead、429/5xx→Transient、400→RequestRejected)。空向量/维度混乱 → ResultInvalidError。
8. 压测 harness(tools/soak/,不被 import、不入 pytest 门)
结构(场景矩阵与不变量以 findings 文档为准,此处只定工程形态):
| 组件 | 职责 |
|---|---|
tools/soak/corpus.py |
语料装载:harness.db traces→messages 还原器、telemetry.db 回放抽取、vt_frames/chs_images 组装 |
tools/soak/scenarios.py |
P1-P6 请求生成器(每场景一个 async 生成器,产出 chat/embed 调用参数) |
tools/soak/run_soak.py |
入口:--scenario P3 --budget-calls N --budget-tokens M;开跑前 FLUSHDB(仅 db3,校验连接串含 /3 否则拒跑)+ cache_namespace=run_id;双上限(请求数/token)任一命中即停 |
tools/soak/scoreboard.py |
跑后读本 run 遥测库断言 6 条硬不变量(findings §4);产报告至 tests/outputs/soak/<run_id>.md |
故障源混编按 findings §3 配置(坏 key/黑洞/紧看门狗/紧闸源);回放请求掺 cache_salt=run_id。harness 用库自身的 GatewayClient/EmbeddingClient 与遥测(吃狗粮),不引入第三方压测框架。
8.1 待人类签字项(harness 实现前拍板,可在批准本设计时一并给)
| 项 | 提议值(可改) |
|---|---|
| 预算上限 | P1≤500 次 / P2≤450 次 / P3+P4≤2500 次 / P5≤1000 次 / P6≤8000 次或 3h;全程 token 硬顶 5000 万(输入+输出合计,按遥测实测累计) |
| 网关保护 | harness 全局闸:max_concurrency=100、全局 RPM=600;跑 P6 建议夜间时段(具体由人类定) |
| P6 混合比例 | P1 10% / P2 20% / P3 50%(其中 3/10 为重复 messages 走缓存双向,即 P4)/ P5 20% |
9. 旧版行为审计(迁移类,逐条标注)
限流(CHS limiter.py + scripts.py):六闸顺序与判据/单 Lua 原子/ZSET 租约双保险/服务器时钟窗口/settle 落 acquire 窗口/负数 INCRBY/幂等 flag/拒绝零副作用/进度键 TTL 3600 —— 全部保留。_BACKOFF_S=0.05 固定轮询的阻塞 acquire —— 替换为契约既有 acquire(M1 形态,poll_interval 可配)。key 前缀 cclimit: —— 替换为 pgw:limit:。LimiterError —— 替换为 GovernanceBackendError(M1 错误模型)。
熔断(CHS provider_gate.py):五操作 Lua 语义/epoch 仅开断 +1/探针租约惰性重发/release_probe 立即可接管/force_open 顶格/时钟回拨 clamp —— 全部保留。probe_ttl 由 container 自动算 —— 替换为显式配置 + 装配守卫(§3)。key 前缀 —— 替换为 pgw:gate:。
背压(CHS governance.py):双条件公式/时钟分工/poll jitter —— 保留;ProviderUnavailableError("stalled") —— 替换为 AllSourcesExhausted(reason="stalled")(M1 错误模型翻译,语义等价)。
Embedding(GovDoc/VT):逐条裁决见 §7.3 表(该表即审计,含三处"有意放弃":GovDoc 自研退避、on_usage 回调、VT 同步接口)。
遥测(GovDoc taskrun 蓝本):asyncpg 用法/占位符/幂等写法 —— 保留;失败冒泡方向 —— 有意反转为静默降级(遥测铁律)。
10. 非功能维度
- 并发与取消:Redis 后端所有 await 点可被取消;单条 Lua 执行期不可中断但均为毫秒级短脚本;permit 取得后的取消由 M1
try/finally结算释放路径覆盖(retry.py:160-170,行为不变,Redis 后端下集成测试重验);acquire/stall 等待的 sleep 可取消;Embedding 循环取消穿透同 RetryMW 契约(探针取消 → release_probe);Postgres 写入被取消 → 该行丢失,可接受(遥测非事务性承诺)。 - 降级方向:限流/熔断 Redis 掉线 → 准入侧
GovernanceBackendError报错不放行(六处 RedisError 捕获点全 fail-closed,CHS 同款);释放侧(settle/release/release_probe)失败 → warning 降级不掩盖主异常(M1 勘误);Postgres 遥测/缓存 → 静默降级;pricing 查不到 → cost=None + warning;pricing 文件坏 → 装配报错。 - 幂等与重复:settle/release/release_probe 幂等(契约已测,Redis 版靠 Lua 天然幂等 + flag);遥测
ON CONFLICT DO NOTHING;harness FLUSHDB 可重复;EmbeddingClient.embed无副作用可重放(网关计费除外)。 - 持久化与原子性:六闸判定与占用同一 Lua 原子;gate 每操作单 Lua 原子;跨 key(全局+单源)在同一脚本内一致;Postgres 单行 INSERT 原子;部分写入不产生半行。Redis 重启 = 治理状态清零(限额窗口/熔断状态重新累积),有意接受(CHS 同款,治理状态是软状态)。
11. 错误分类与测试策略
分类归属:Redis 不可用(准入侧)→ GovernanceBackendError;stall → AllSourcesExhausted("stalled");embedding 网关错误 → 四分类同款翻译;维度不符/空向量 → ResultInvalidError(不熔断);pricing 缺价 → 非错误(cost=None)。
测试(全部先红后绿):
- 契约:conftest 增
redisparam(缺REDIS_URLskip;唯一 scope 隔离;db3),两契约文件零改动全绿;时间语义补小 TTL 真实等待的 integration 变体。 - 保真:Lua 移植逐段比对 CHS 行号(plan 内设检查点);移植 CHS 三个跨连接集成用例(双连接池共享全局并发/全局 RPM/进度可见)。
- 联合验证(③):多 client 双连接池下,多源换源 + 全局 RPM 不超配 + 熔断状态跨连接共享的集成测试;取消穿透在 Redis 后端下重验(in-flight 取消 → lease 释放)。
- stall:memory 后端 FakeClock 推进双条件各自与同时成立的四象限;fail_fast 路径不受影响回归。
- Postgres:真实实验室 Postgres?——无现成实例,用 conda 环境本地起临时 postgres 或 docker;若都不可用则该集成测试标 skip 并在验收注明(integration;幂等/降级/并发 50 写并落全部)。
- Embedding:unit(ScriptedTransport 分批/保序/归一化/维度校验/重试换源)+ e2e 真实网关冒烟(若网关有 embedding 端点;没有则 e2e 降为对 MiniMax chat 网关的 404 行为记录,unit 全覆盖)。
- harness:本体不入 pytest;
corpus.py的 traces 还原器与不变量断言函数给 unit 测试(纯函数)。
覆盖率目标沿用 ≥80%;遥测/缓存降级、限流结算退款、熔断开路半开、Redis 掉线方向仍是一等测试对象。
12. 内部顺序与交付物
沿 ROADMAP §3:① Redis 限流 → ② Redis 熔断 → ③ 联合验证 → ④ stall → ⑤ Postgres 遥测 + ⑥ pricing(可并行)→ ⑦ Embedding → ⑧ harness(⑤-⑧ 相互独立,⑧ 依赖全部)。交付物:backends/redis/{limiter,breaker}.py、middleware/retry.py stall 增量、telemetry/postgres.py、pricing.py、embedding.py + transport 增量 + ports.py/types.py 新增(EmbeddingTransport/EmbeddingResponse,只增不改)、config 增量、tools/soak/、.env.example 回填、迁移文档 embedding 条目更新。
13. 开放问题(设计内已给提议,批准时可一并裁决)
- §8.1 三项签字(预算/网关保护/P6 比例)。
- Postgres 集成测试环境:实验室有无可用 Postgres 实例?(无则按 §11.5 降级方案)
- 真实网关是否有 embedding 端点可供 e2e?(无则按 §11.6 降级方案)