From 724dc3c3286f087ae4d97aeba93d0a38ec0e138d Mon Sep 17 00:00:00 2001 From: iomgaa Date: Tue, 21 Jul 2026 01:28:33 -0400 Subject: [PATCH] docs: backfill env template and migration notes for M2 --- .env.example | 43 ++++++++++++++++--- research-wiki/migrations/govdoc-saas.md | 2 +- research-wiki/migrations/video-tree-trm5.md | 4 +- src/polygateway/embedding.py | 8 ++-- .../integration/test_redis_governance_time.py | 4 +- tests/unit/test_embedding.py | 8 +--- tests/unit/test_soak_corpus.py | 7 ++- tools/soak/run_soak.py | 4 +- tools/soak/scenarios.py | 4 +- tools/soak/scoreboard.py | 5 +-- 10 files changed, 59 insertions(+), 30 deletions(-) diff --git a/.env.example b/.env.example index b2d64b6..8008a67 100644 --- a/.env.example +++ b/.env.example @@ -33,22 +33,51 @@ LLM_CIRCUIT_BREAKER_COOLDOWN=60 # 或 LLM__BREAKER__COOLDOWN_S # LLM_TIMEOUT=120 # 源缺 TIMEOUT_S 时的缺省 # LLM_TTFT_TIMEOUT=30 # 平铺看门狗缺省(成对生效) # LLM_INTER_TOKEN_TIMEOUT=15 -# LLM__BREAKER__PROBE_TTL_S=240 # 缺省派生: max(2×最大源超时, cooldown) -# LLM__BACKPRESSURE__STALL_WINDOW_S=300 # M2 启用 stall 判定 +# LLM__BREAKER__PROBE_TTL_S=240 # 缺省派生: max(2×最大源超时, cooldown, 最大源超时+5);显式值须 ≥ 最大源超时+5 +# LLM__BACKPRESSURE__STALL_WINDOW_S=300 # stall 双条件判死窗口;须 ≥ 最大源 TTFT # LLM__BACKPRESSURE__POLL_INTERVAL_S=0.05 # LLM__SELECTOR=round_robin # round_robin(默认) | least_inflight # LLM__QUOTA_FULL=wait # wait(默认) | fail_fast # ══ 装配选择(PGW_*)══ -PGW_LIMITER_BACKEND=memory # M1 仅 memory;redis 随 M2 -PGW_BREAKER_BACKEND=memory +PGW_LIMITER_BACKEND=memory # memory | redis(redis 需 REDIS_URL;多进程 worker 必须 redis) +PGW_BREAKER_BACKEND=memory # memory | redis PGW_CACHE_BACKEND=none # redis | memory | none(必填,显式优于隐式) -PGW_TELEMETRY_BACKEND=none # sqlite | none(postgres 随 M2) +PGW_TELEMETRY_BACKEND=none # sqlite | postgres | none(必填) # PGW_TELEMETRY_SQLITE_PATH=logs/telemetry.db # sqlite 时必填 +# PGW_TELEMETRY_PG_DSN=postgresql://user:pass@host:5432/polygateway # postgres 时必填;严禁指向在用业务库(实验室约定: 专用库 polygateway) +# PGW_PRICING_PATH=config/prices.json # 可选: {"": {"input_per_1m": x, "output_per_1m": y}};缺省 cost 恒 None # PGW_CACHE_NAMESPACE=<项目名或租户前缀> # 缓存启用时必填(防跨项目毒化) # PGW_CACHE_TTL_S=604800 # 缓存启用时必填,须 > 0 # PGW_STRUCTURED_MAX_RETRIES=1 # 0 = 解析失败不重问(CHS 策略) # PGW_LEASE_TTL_S=1500 # permit 租约;须 ≥ 最大源 timeout -# ══ Redis(缓存;M2 起亦供分布式限流/熔断)══ -# REDIS_URL=redis://localhost:6379/0 +# ══ Redis(缓存 + 分布式限流/熔断)══ +# 实验室纪律: 共享实例的 db0 有在用键,PolyGateway 一律用专用 db3(soak 会 FLUSHDB!) +# REDIS_URL=redis://:password@host:6379/3 + +# ══ Embedding scope(M2;EmbeddingClient.from_env 装配)══ +# EMBED__QWEN__1__BASE_URL= +# EMBED__QWEN__1__API_KEY= +# EMBED__QWEN__1__MODEL=text-embedding-v3 +# EMBED__QWEN__1__TIMEOUT_S=60 +# EMBED__RETRY__MAX_ATTEMPTS=3 +# EMBED__RETRY__BACKOFF_BASE_S=1.0 +# EMBED__RETRY__BACKOFF_MAX_S=10.0 +# EMBED__BREAKER__FAIL_THRESHOLD=5 +# EMBED__BREAKER__COOLDOWN_S=60 +# EMBED__BATCH_SIZE=64 # 必填: 每批条数(分批是行为关键,不设默认) +# EMBED__NORMALIZE=false # 可选: true = 库内 L2 归一化(VT 语义) +# EMBED__EXPECTED_DIM=768 # 可选: 维度校验,不符抛 ResultInvalid + +# ══ SOAK scope(压测 harness 专用;tools/soak/run_soak.py --scope SOAK)══ +# 健康源 + 故障源混编(findings §3): 坏 key 源 / 黑洞源 / 紧看门狗源 / 紧闸源 +# SOAK__MINIMAX__1__BASE_URL= # 健康源(真实网关) +# SOAK__MINIMAX__1__API_KEY= +# SOAK__MINIMAX__1__MODEL= +# SOAK__MINIMAX__1__TIMEOUT_S=180 +# SOAK__MINIMAX__2__BASE_URL= # 坏凭据源: 真实网关 + 错误 API_KEY → 401 +# SOAK__MINIMAX__2__API_KEY=sk-wrong-key-on-purpose +# SOAK__MINIMAX__3__BASE_URL=http://10.255.255.1/v1 # 黑洞源: 防火墙 DROP → 连接超时 +# SOAK__MINIMAX__4__RPM=5 # 紧闸源: 真实源 + rpm=5 / 并发 1 +# SOAK__MINIMAX__4__MAX_CONCURRENCY=1 diff --git a/research-wiki/migrations/govdoc-saas.md b/research-wiki/migrations/govdoc-saas.md index 69669b9..970264f 100644 --- a/research-wiki/migrations/govdoc-saas.md +++ b/research-wiki/migrations/govdoc-saas.md @@ -32,7 +32,7 @@ | `.../llm/__init__.py` | 0 | 空 | **删除**(随子包) | | `.../docagent_core/protocols.py` | 151 | 10 个共享 Protocol,其中 LLM 相关 2 个 | **改写**: 删 `TelemetryRecorder`(唯一消费者是被删的 client.py/telemetry_sqlite.py);`LLMProvider` 保留(agent/workflow 的消费契约,GatewayClient 结构化满足);其余 8 个非 LLM 端口不动 | | `.../docagent_core/types.py` | 25 | `LLMResponse` frozen dataclass(11 字段) | **改写**: 改为 re-export 库的 `LLMResponse`(一行 shim,超集兼容),下游 import 路径零改动 | -| `.../retrieval/embedding.py` | 165 | Embedding 客户端,含独立手写重试(embedding.py:107-131,第三处重复) | **保留(暂)**: 待 ARCHITECTURE Q3 决策(建议 M2 纳库),本次迁移不动 | +| `.../retrieval/embedding.py` | 165 | Embedding 客户端,含独立手写重试(embedding.py:107-131,第三处重复) | **删除,换 `polygateway.EmbeddingClient`**(Q3 已拍板,M2 已交付)。映射: `batch_size` → `EMBED__BATCH_SIZE`;`dimension` 校验 → `EMBED__EXPECTED_DIM`(不符抛 ResultInvalidError,原 EmbeddingUnavailableError 语义由四分类承接);自研退避 → 库退避(封顶+jitter,原版无上限无 jitter 为有意升级);`on_usage` 回调 → 遥测 llm_calls 行(prompt_tokens)+ pricing cost,消费方改读遥测库 | | `packages/docagent-core/tests/test_breaker.py` | 32 | 熔断状态机单测 | **删除**(等价测试随库实现交付) | | `packages/docagent-core/tests/test_imports.py` | 10 | 公共 API import 冒烟(第 5-6 行 import 被删模块) | **改写** import 目标 | | `packages/docagent-core/tests/conftest.py` | — | `ScriptedLLM` fake(构造 `LLMResponse`) | **保留**(经 types.py re-export 零改动) | diff --git a/research-wiki/migrations/video-tree-trm5.md b/research-wiki/migrations/video-tree-trm5.md index bf0af32..0e539d1 100644 --- a/research-wiki/migrations/video-tree-trm5.md +++ b/research-wiki/migrations/video-tree-trm5.md @@ -28,7 +28,7 @@ | `adapters/telemetry.py` | 229 | SQLite 遥测(单连接+锁+WAL+幂等+to_thread+降级) | **删除**(库 `telemetry/sqlite.py` 继任) | | `adapters/ocr.py` | 128 | MonkeyOCR `/ocr/text` **裸调**:同步 requests、线程轮询双端点、单帧失败跳过 | **删除**(换库 `OcrTextPort`,升级为全治理;行过滤/拼接逻辑上移业务侧,见 §4) | | `adapters/vlm.py` | 131 | base64 编码图片并注入最后一条 user message,委托 LLM client | **保留改写**(业务侧封装:`_inject_images` 保留,委托对象换成库 client) | -| `adapters/embedding.py` | 184 | local(sentence-transformers)/remote(OpenAI SDK) 嵌入 | **保留**(Q3 未决;remote 实测**无任何重试**且为同步调用,见 §9-R11) | +| `adapters/embedding.py` | 184 | local(sentence-transformers)/remote(OpenAI SDK) 嵌入 | **remote 路径删除,换 `polygateway.EmbeddingClient`**(Q3 已拍板纳入 M2,2026-07-20;库已交付);local(sentence-transformers)保留。迁移要点: 库为 async `embed(list[str]) -> EmbeddingResponse(list[list[float]])`——VT 同步调用点需自包同步壳(`asyncio.run` 或事件循环内 await),ndarray/tensor 转换自带 adapter(`np.asarray(resp.vectors, dtype=np.float32)`);装配传 `EMBED__NORMALIZE=true` 保 L2 归一化语义 | | `adapters/baseline_diagnosis_store.py` | 183 | 基线诊断结果 SQLite 存储(业务数据) | **保留**(业务侧) | | `core/protocols.py` | 69 | `LLMProvider`/`VLMProvider`/`TelemetryRecorder` 三 Protocol | **改写**(LLMProvider 由库满足;VLMProvider 指向改写后的 vlm.py;TelemetryRecorder 删除,见 §3) | | `core/types.py` | 182 | `LLMResponse`(11 字段) + 业务类型 | **改写**(LLMResponse 改为 re-export 库类型;其余保留) | @@ -49,7 +49,7 @@ | `LLMResponse`(11 字段) (`core/types.py:19-29`) | 库 `LLMResponse`(§5.1 超集) | 字段逐一比对**完全一致**(content/thinking/model/provider/prompt_tokens/completion_tokens/latency_ms/ttft_ms/max_inter_token_ms/cache_hit/call_id);库新增 source_name/cost/usage_source 只增不删 | `core/types.py` 改 re-export:`from polygateway import LLMResponse` | | `CircuitOpenError` (`adapters/llm.py:36`) | 库 `CircuitOpenError`(§6.1) | 业务侧无 except 该异常(grep 实测),仅治理层内部 | 无需 shim | | `MonkeyOCRClient.transcribe_frames(frame_paths) -> str` (`adapters/ocr.py:88`) | `OcrTextPort.recognize_text(bytes)`(§7.10) | **不兼容**:项目端口收路径列表、返回拼接文本、单帧失败跳过;库端口收单帧 bytes、失败抛异常 | 业务适配器(约 30 行):读文件→逐帧 `recognize_text`→行过滤去重→`"帧N: ..."` 拼接→单帧异常捕获跳过;实现 `app/ports.py:OCRProvider` 不变,`vision.py:105` 调用点零改动 | -| `EmbeddingProvider`(同步 `embed()`) (`app/ports.py:17-36`) | 暂无(Q3 开放) | 若 M2 纳入:库必为异步端口,同步调用点需适配 | 本次迁移不动 | +| `EmbeddingProvider`(同步 `embed()`) (`app/ports.py:17-36`) | `polygateway.EmbeddingClient.embed`(async,M2 已交付) | 同步→异步适配 + list[list[float]]→ndarray 转换留业务侧 adapter;normalize 走库开关 | M4 执行 | --- diff --git a/src/polygateway/embedding.py b/src/polygateway/embedding.py index ffd04c6..78d32ed 100644 --- a/src/polygateway/embedding.py +++ b/src/polygateway/embedding.py @@ -158,7 +158,9 @@ class EmbeddingClient: outcomes = [] for start in range(0, len(texts), self._batch_size): outcomes.append( - await self._embed_batch(texts[start : start + self._batch_size], session_id, parent_call_id) + await self._embed_batch( + texts[start : start + self._batch_size], session_id, parent_call_id + ) ) return self._merge(outcomes) @@ -394,9 +396,7 @@ class EmbeddingClient: def _total_cost(self, outcomes: list[_BatchOutcome]) -> float | None: if self._pricing is None: return None - costs = [ - self._pricing.cost(o.source.model, o.result.prompt_tokens, 0) for o in outcomes - ] + costs = [self._pricing.cost(o.source.model, o.result.prompt_tokens, 0) for o in outcomes] known = [c for c in costs if c is not None] return sum(known) if known else None diff --git a/tests/integration/test_redis_governance_time.py b/tests/integration/test_redis_governance_time.py index 87520b2..dd14471 100644 --- a/tests/integration/test_redis_governance_time.py +++ b/tests/integration/test_redis_governance_time.py @@ -56,9 +56,7 @@ def _advance_dependent_cases() -> set[str]: def test_meta_variants_cover_all_time_cases(): """完整性守卫: 契约的时间用例集合 == 本文件变体集合(1:1 映射表机械化)。""" variants = { - name.removeprefix("test_variant_") - for name in globals() - if name.startswith("test_variant_") + name.removeprefix("test_variant_") for name in globals() if name.startswith("test_variant_") } assert variants == _advance_dependent_cases() diff --git a/tests/unit/test_embedding.py b/tests/unit/test_embedding.py index a2ed5d9..65e548a 100644 --- a/tests/unit/test_embedding.py +++ b/tests/unit/test_embedding.py @@ -94,14 +94,10 @@ class TestEmbedTransport: assert payload == {"model": "embed-1", "input": ["a", "b"]} return httpx.Response( 200, - json=_ok_body( - [[1.0, 0.0], [0.0, 1.0]], usage={"prompt_tokens": 5}, shuffle=True - ), + json=_ok_body([[1.0, 0.0], [0.0, 1.0]], usage={"prompt_tokens": 5}, shuffle=True), ) - result = await _transport_with(handler).embed( - texts=["a", "b"], source=_src(), call_id="c" - ) + result = await _transport_with(handler).embed(texts=["a", "b"], source=_src(), call_id="c") assert result.vectors == [[1.0, 0.0], [0.0, 1.0]] # 乱序响应按 index 重排 assert result.dim == 2 assert result.prompt_tokens == 5 and result.usage_source == "measured" diff --git a/tests/unit/test_soak_corpus.py b/tests/unit/test_soak_corpus.py index ec9f792..e7f1b57 100644 --- a/tests/unit/test_soak_corpus.py +++ b/tests/unit/test_soak_corpus.py @@ -28,8 +28,11 @@ def _mini_harness_db(path): " completion_tokens INTEGER, stop_reason TEXT, steps_json TEXT)" ) steps = [ - {"thought": f"思考{i}", "tool_call": {"tool": "view_node", "args": {"i": i}}, - "tool_output": f"输出{i}" * 50} + { + "thought": f"思考{i}", + "tool_call": {"tool": "view_node", "args": {"i": i}}, + "tool_output": f"输出{i}" * 50, + } for i in range(3) ] conn.execute( diff --git a/tools/soak/run_soak.py b/tools/soak/run_soak.py index 308c57c..ee95913 100644 --- a/tools/soak/run_soak.py +++ b/tools/soak/run_soak.py @@ -145,7 +145,9 @@ def _scoreboard(args: argparse.Namespace, env: dict[str, str]) -> None: results.append(json.loads(path.read_text(encoding="utf-8"))) rows = sb.load_rows(*(r["telemetry"] for r in results)) calls = sum(r["stats"]["calls"] for r in results) - max_attempts = int(env.get(f"{args.scope}__RETRY__MAX_ATTEMPTS", env.get("LLM_MAX_RETRIES", "3"))) + max_attempts = int( + env.get(f"{args.scope}__RETRY__MAX_ATTEMPTS", env.get("LLM_MAX_RETRIES", "3")) + ) verdicts: list[tuple[str, str]] = [] def _check(name: str, fn, *fargs, **fkwargs) -> None: diff --git a/tools/soak/scenarios.py b/tools/soak/scenarios.py index 1642b51..7b32df3 100644 --- a/tools/soak/scenarios.py +++ b/tools/soak/scenarios.py @@ -72,7 +72,9 @@ def weighted_mix(weights: dict[str, float], rng) -> str: return next(reversed(weights)) -async def p1_trace_chains(corpus: SoakCorpus, run_id: str, rng=random.random) -> AsyncIterator[Item]: +async def p1_trace_chains( + corpus: SoakCorpus, run_id: str, rng=random.random +) -> AsyncIterator[Item]: """P1 文本长上下文回放: 单链串行,session/parent 链路照原样语义。""" for chain_idx, chain in enumerate(corpus.chains): session_id = f"{run_id}-p1-{chain_idx}" diff --git a/tools/soak/scoreboard.py b/tools/soak/scoreboard.py index 14e4ae4..220b81b 100644 --- a/tools/soak/scoreboard.py +++ b/tools/soak/scoreboard.py @@ -48,6 +48,7 @@ def inv_call_ids_unique(rows: list[Row]) -> None: def _minute_bucket(created_at: str) -> int: return int(datetime.fromisoformat(str(created_at)).timestamp()) // 60 + def inv_rpm_never_exceeded(rows: list[Row], per_source_rpm: dict[str, int]) -> None: """不变量 3: 按遥测时间戳重算,任一分钟桶内单源请求数 ≤ RPM 配置。 @@ -59,9 +60,7 @@ def inv_rpm_never_exceeded(rows: list[Row], per_source_rpm: dict[str, int]) -> N for r in rows if per_source_rpm.get(r["source_name"], 0) > 0 ) - breaches = { - key: n for key, n in buckets.items() if n > per_source_rpm[key[0]] - } + breaches = {key: n for key, n in buckets.items() if n > per_source_rpm[key[0]]} assert not breaches, f"RPM 击穿: {dict(list(breaches.items())[:5])}"