"""M5 端到端:群像生成→入库闭环 + 切 provider 回退/降级回归(真实 DB,零 token)。 证明 M5 闭环(DoD:世界观/群像生成入库过 continuity 校验;多 provider 回退/降级——切 provider 不破): 用例 1 · 群像生成 → 409 gate → 确认 → 入库(DB 真源 + schema→JSONB 形变) `POST /world/generate`(writer 假网关产 `WorldGenResult`)→ 预览实体(不入库)。 `POST /characters/generate`(writer 假网关产 `CharacterGenResult`,群像 3 卡)→ 预览。 `POST /characters`(入库)双 gate: (a) precheck(analyst 假网关产带冲突 `ContinuityReview`)有冲突且 `acknowledge_conflicts` 缺省 → **409 CONFLICT_UNRESOLVED** + `details.conflicts`(不静默入库,守不变量 #3)。 (b) `acknowledge_conflicts=true` 重发 → 201;断言 `characters` 行**真落 pg**,且 schema list/str → DB JSONB dict 形变(traits/speech_tics→`{"items":[...]}`、 arc→`{"text":...}`)。 入库网关按 `req.output_schema` 路由:CharacterGenResult / WorldGenResult / ContinuityReview。 用例 2 · 越权写表被白名单丢弃(partition_writes,§5.6) character-gen 只声明 `writes=["characters"]`;构造一张越权写表只能经端点 `partition_writes` 丢弃——这里证明合法路径(characters)放行、`rejected_tables` 为空(正常态)。 用例 3 · 切 provider 回退(primary transient 失败 → fallback 服务;DB 记账记实际服务方) 端到端 HTTP(角色生成)跑在**真实 `Gateway`**(非 stub)上,链 = primary(deepseek) → fallback(openai)。primary 假适配器每次 `complete` 抛 `TransientProviderError`(耗尽重试) → 网关沿链切 fallback 服务 → 端点 201。`served_by` 不出 API,故经 **DB 真源 `usage_ledger`** 断言实际服务方是 **openai(fallback)而非 deepseek(primary)**——证明回退真发生、记账记实际 服务方(ARCH §4.5)。 用例 4 · 能力降级(primary 不支持结构化输出 → 链上有支持者则重排/降级,§4.4) primary(deepseek) `capabilities().structured_output=False` 且无 fallback → 网关降级用它 (适配器自走 JSON 提示返 parsed),`served_by.degraded=True`。经真 `Gateway.run` 直测 (聚焦集成测,证明降级路径真走,不经 HTTP 更清晰)。 确定性 & 零成本(同 M1–M4):真实 `Gateway` + 多 provider 假适配器(据 `req.output_schema` 分支返回固定 `parsed`,绝不联网)+ 真实 `SqlAlchemyLedgerSink`(记账真落 pg)。无 pg → skip。 坑(见 memory/gotchas): - 生成网关是**请求 scope 依赖**(`get_worldbuilder_gateway`/`get_character_gen_gateway`/ `get_precheck_gateway`),非 background 自建——直接 `app.dependency_overrides` 注假网关即可, 无需 monkeypatch `build_gateway_for_tier`(与 M4 学文风后台自建网关不同)。 - 真 `Gateway.run` 据 `req.output_schema` 选 schema:CharacterGenResult/WorldGenResult/ ContinuityReview——假适配器照此分支。 - `get_sessionmaker` engine 绑首个 loop → 每 DB 测试清缓存重建、结束 dispose(同 M3/M4)。 - ASGITransport 默认 `raise_app_exceptions=True`;snake_case; 409 envelope = `{error:{code,message,details}}`。 - usage_ledger 记账记**实际服务方**:回退后 provider 列是 fallback 名,非 primary——这是 「切 provider 不破」的 DB 真源证据(served_by 不出 API)。 """ from __future__ import annotations import uuid from collections.abc import AsyncIterator, Callable from typing import Annotated import httpx import pytest from asgi_lifespan import LifespanManager from fastapi import Depends from sqlalchemy import delete, select from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker from ww_agents import ( CharacterCard, CharacterGenResult, Conflict, ContinuityReview, WorldEntityCard, WorldGenResult, ) from ww_db import get_session, get_sessionmaker from ww_db.models import Character, Project, UsageLedger from ww_llm_gateway import ( Gateway, SqlAlchemyLedgerSink, chain_from_routing, resolve_route, ) from ww_llm_gateway.adapters.base import ( Capabilities, ProviderResult, ProviderUsage, StreamChunk, ) from ww_llm_gateway.errors import TransientProviderError from ww_llm_gateway.routing import Route from ww_llm_gateway.types import LlmRequest, Scope, Tier # writer/analyst 默认都路由到 deepseek(config.tier_defaults);用例 3 显式构 deepseek→openai 链。 _PRIMARY = "deepseek" _FALLBACK = "openai" # 各档位假用量(喂记账;证明记账落库 + 实际服务方归因)。 _USAGE = { "writer": ProviderUsage(input_tokens=31, output_tokens=17), "analyst": ProviderUsage(input_tokens=23, output_tokens=7), } # 群像假产物:3 张差异化角色卡(schema list/str 形,入库时由写侧 repo 转 DB JSONB dict)。 _CHARACTER_CARDS = CharacterGenResult( cards=[ CharacterCard( name="叶无尘", role="主角", traits=["隐忍", "腹黑"], backstory="灭门遗孤,背负血仇。", arc="从复仇者走向守护者。", speech_tics=["呵,有意思", "不过如此"], tags=["扮猪吃虎", "天才"], relations=[], ), CharacterCard( name="苏清歌", role="女主", traits=["果决", "外冷内热"], backstory="圣女出身,叛出师门。", arc="从教条走向自我。", speech_tics=["哼"], tags=["大小姐", "御姐"], relations=[], ), CharacterCard( name="胖三", role="工具人", traits=["市侩", "重情"], backstory="市井小贩,江湖百晓生。", arc="从贪财到舍命相护。", speech_tics=["客官您内行"], tags=["搞笑担当"], relations=[], ), ] ) # 世界观假产物:两个硬规则实体(预览,不入库)。 _WORLD_ENTITIES = WorldGenResult( entities=[ WorldEntityCard(type="力量体系", name="九转玄功", rules=["每转需一枚玄晶", "不可逆修"]), WorldEntityCard(type="势力", name="天机阁", rules=["阁主不入世俗纷争"]), ] ) # precheck 假冲突(让入库 gate 触发 409)。 _PRECHECK_CONFLICT = ContinuityReview( conflicts=[ Conflict( type="设定违例", where="叶无尘", refs=["九转玄功"], suggestion="叶无尘的功法与九转玄功『不可逆修』硬规则冲突,请调整设定", ) ] ) # precheck 无冲突(acknowledge 后重发——实际仍会返冲突,但带 ack 放行)。 class _FakeGenAdapter: """实现 `ProviderAdapter` Protocol:按 `req.output_schema` 分支返回 parsed,绝不联网。 - `CharacterGenResult` → 3 张群像卡(writer)。 - `WorldGenResult` → 两个世界观实体(writer)。 - `ContinuityReview` → 带冲突的预检结果(analyst)——触发入库 409 gate。 每路带 `ProviderUsage` → 网关据此落 usage_ledger。 """ def __init__(self, provider: str = _PRIMARY, *, structured: bool = True) -> None: self.provider = provider self._structured = structured def capabilities(self) -> Capabilities: return Capabilities(structured_output=self._structured) async def complete(self, req: LlmRequest, model: str) -> ProviderResult: schema = req.output_schema if schema is CharacterGenResult: return ProviderResult( text=_CHARACTER_CARDS.model_dump_json(), parsed=_CHARACTER_CARDS, usage=_USAGE["writer"], ) if schema is WorldGenResult: return ProviderResult( text=_WORLD_ENTITIES.model_dump_json(), parsed=_WORLD_ENTITIES, usage=_USAGE["writer"], ) if schema is ContinuityReview: return ProviderResult( text=_PRECHECK_CONFLICT.model_dump_json(), parsed=_PRECHECK_CONFLICT, usage=_USAGE["analyst"], ) raise AssertionError(f"unexpected output_schema in M5 fake adapter: {schema!r}") async def stream(self, req: LlmRequest, model: str) -> AsyncIterator[StreamChunk]: # 生成走 run()(结构化),不走 stream——但 Protocol 要求实现。 yield StreamChunk(usage=_USAGE["writer"]) raise AssertionError("M5 generation must not stream") class _FailingPrimaryAdapter: """primary 假适配器:每次 `complete` 抛 `TransientProviderError`(耗尽重试 → 切回退)。""" provider = _PRIMARY def __init__(self) -> None: self.complete_calls = 0 def capabilities(self) -> Capabilities: return Capabilities(structured_output=True) async def complete(self, req: LlmRequest, model: str) -> ProviderResult: self.complete_calls += 1 raise TransientProviderError("simulated 503", provider=self.provider) async def stream(self, req: LlmRequest, model: str) -> AsyncIterator[StreamChunk]: raise TransientProviderError("simulated 503", provider=self.provider) yield StreamChunk() # pragma: no cover @pytest.fixture async def e2e_sm() -> AsyncIterator[async_sessionmaker[AsyncSession]]: """真实 DB session 工厂;无 pg 时跳过(每测试清缓存重建 engine、结束 dispose)。""" get_sessionmaker.cache_clear() maker = get_sessionmaker() try: async with maker() as probe: await probe.execute(select(1)) except Exception: pytest.skip("postgres not reachable") yield maker await maker.kw["bind"].dispose() get_sessionmaker.cache_clear() def _single_provider_override( adapter: _FakeGenAdapter, ) -> Callable[[AsyncSession], Gateway]: """依赖覆盖:真实 Gateway + 单 provider 假适配器 + 真实 ledger(请求 session)。""" def _override(session: Annotated[AsyncSession, Depends(get_session)]) -> Gateway: return Gateway( adapters={adapter.provider: adapter}, ledger=SqlAlchemyLedgerSink(session), resolver=resolve_route, ) return _override def _fallback_chain_override( primary: _FailingPrimaryAdapter, fallback: _FakeGenAdapter, ) -> Callable[[AsyncSession], Gateway]: """依赖覆盖:真实 Gateway + deepseek→openai 回退链 + 真实 ledger。 用真 `Gateway`(非 stub)走 §4.5 回退逻辑:primary transient 失败耗尽重试 → 切 fallback。 链经 `chain_from_routing` 构造(同 apps/api build_gateway_for_tier 用法)。 """ def _resolver(_tier: Tier) -> list[Route]: return chain_from_routing(_tier, f"{_PRIMARY}:deepseek-chat", [f"{_FALLBACK}:gpt-4o"]) def _override(session: Annotated[AsyncSession, Depends(get_session)]) -> Gateway: return Gateway( adapters={_PRIMARY: primary, _FALLBACK: fallback}, ledger=SqlAlchemyLedgerSink(session), chain_resolver=_resolver, max_retries=1, ) return _override async def _cleanup(e2e_sm: async_sessionmaker[AsyncSession], project_uuid: uuid.UUID) -> None: """按 FK 顺序清理。""" async with e2e_sm() as cleanup: await cleanup.execute(delete(UsageLedger).where(UsageLedger.project_id == project_uuid)) await cleanup.execute(delete(Character).where(Character.project_id == project_uuid)) await cleanup.execute(delete(Project).where(Project.id == project_uuid)) await cleanup.commit() async def test_m5_generate_characters_gate_acknowledge_and_persist( e2e_sm: async_sessionmaker[AsyncSession], ) -> None: """用例 1+2:世界观/群像生成预览 → 入库 409 gate → 确认放行 → 真落 pg(schema→JSONB 形变)。""" from ww_api.main import create_app from ww_api.services.project_deps import ( get_character_gen_gateway, get_precheck_gateway, get_worldbuilder_gateway, ) writer_adapter = _FakeGenAdapter(provider=_PRIMARY) analyst_adapter = _FakeGenAdapter(provider=_PRIMARY) app = create_app() # 生成网关是请求 scope 依赖 → 直接 override 注假网关(无需 monkeypatch build_gateway)。 app.dependency_overrides[get_worldbuilder_gateway] = _single_provider_override(writer_adapter) app.dependency_overrides[get_character_gen_gateway] = _single_provider_override(writer_adapter) app.dependency_overrides[get_precheck_gateway] = _single_provider_override(analyst_adapter) transport = httpx.ASGITransport(app=app) project_uuid: uuid.UUID | None = None try: async with LifespanManager(app): async with httpx.AsyncClient(transport=transport, base_url="http://test") as client: create_resp = await client.post( "/projects", json={"title": "M5 群像验证作品", "genre": "玄幻"} ) assert create_resp.status_code == 201 project_id = create_resp.json()["id"] project_uuid = uuid.UUID(project_id) # 1) 世界观生成预览(不入库)。 world_resp = await client.post( f"/projects/{project_id}/world/generate", json={"brief": "构建一个修真世界的力量体系与势力"}, ) assert world_resp.status_code == 200 entities = world_resp.json()["entities"] assert [e["name"] for e in entities] == ["九转玄功", "天机阁"] assert entities[0]["rules"] == ["每转需一枚玄晶", "不可逆修"] # 2) 群像生成预览(3 卡,群像防雷同;不入库)。 gen_resp = await client.post( f"/projects/{project_id}/characters/generate", json={"brief": "主角团三人,定位差异化", "count": 3}, ) assert gen_resp.status_code == 200 cards = gen_resp.json()["cards"] assert [c["name"] for c in cards] == ["叶无尘", "苏清歌", "胖三"] assert cards[0]["traits"] == ["隐忍", "腹黑"] # 3a) 入库(未 acknowledge)→ precheck 有冲突 → 409 CONFLICT_UNRESOLVED。 ingest_blocked = await client.post( f"/projects/{project_id}/characters", json={"cards": cards}, ) assert ingest_blocked.status_code == 409 err = ingest_blocked.json()["error"] assert err["code"] == "CONFLICT_UNRESOLVED" details = err["details"] assert details["conflict_count"] == 1 assert details["conflicts"][0]["type"] == "设定违例" assert details["conflicts"][0]["where"] == "叶无尘" # 入库被拦 → 此刻 characters 表应无行(不静默入库,守不变量 #3)。 async with e2e_sm() as mid: pending = ( ( await mid.execute( select(Character).where(Character.project_id == project_uuid) ) ) .scalars() .all() ) assert pending == [] # 3b) 作者裁决后带 acknowledge_conflicts=true 重发 → 201 入库。 ingest_ok = await client.post( f"/projects/{project_id}/characters", json={"cards": cards, "acknowledge_conflicts": True}, ) assert ingest_ok.status_code == 201 body = ingest_ok.json() assert body["created"] == ["叶无尘", "苏清歌", "胖三"] # 用例 2:合法路径(writes=["characters"])无越权 → rejected_tables 空(正常态)。 assert body["rejected_tables"] == [] # DB 真源断言:characters 行真落 pg + schema→DB JSONB dict 形变。 assert project_uuid is not None async with e2e_sm() as verify: rows = ( ( await verify.execute( select(Character) .where(Character.project_id == project_uuid) .order_by(Character.name) ) ) .scalars() .all() ) assert {r.name for r in rows} == {"叶无尘", "苏清歌", "胖三"} ye = next(r for r in rows if r.name == "叶无尘") # 形变:traits/speech_tics list → {"items":[...]};arc str → {"text":...}。 assert ye.traits == {"items": ["隐忍", "腹黑"]} assert ye.speech_tics == {"items": ["呵,有意思", "不过如此"]} assert ye.arc == {"text": "从复仇者走向守护者。"} # tags 是 JSONB list(不形变);relations 是 JSONB list。 assert ye.tags == ["扮猪吃虎", "天才"] assert ye.relations == [] assert ye.role == "主角" assert ye.backstory == "灭门遗孤,背负血仇。" # 记账真落 pg:world generate(writer) + characters generate(writer) + # precheck(analyst, 409 路径) + precheck(analyst, 入库路径) = 4 条。 ledger_rows = ( ( await verify.execute( select(UsageLedger).where(UsageLedger.project_id == project_uuid) ) ) .scalars() .all() ) assert len(ledger_rows) == 4 assert all(r.provider == _PRIMARY for r in ledger_rows) finally: if project_uuid is not None: await _cleanup(e2e_sm, project_uuid) async def test_m5_provider_fallback_serves_when_primary_fails( e2e_sm: async_sessionmaker[AsyncSession], ) -> None: """用例 3:HTTP 生成 primary transient 失败 → fallback 服务;DB 记账记实际服务方(openai)。""" from ww_api.main import create_app from ww_api.services.project_deps import get_character_gen_gateway primary = _FailingPrimaryAdapter() fallback = _FakeGenAdapter(provider=_FALLBACK) app = create_app() # 角色生成网关 = 真 Gateway + deepseek(失败)→openai(成功) 回退链。 # character generate 端点不调 precheck(不入库),故只 override 角色生成网关。 app.dependency_overrides[get_character_gen_gateway] = _fallback_chain_override( primary, fallback ) transport = httpx.ASGITransport(app=app) project_uuid: uuid.UUID | None = None try: async with LifespanManager(app): async with httpx.AsyncClient(transport=transport, base_url="http://test") as client: create_resp = await client.post("/projects", json={"title": "M5 切 provider 作品"}) assert create_resp.status_code == 201 project_id = create_resp.json()["id"] project_uuid = uuid.UUID(project_id) # primary transient 失败 → 网关沿链切 fallback → 仍 200(切 provider 不破)。 gen_resp = await client.post( f"/projects/{project_id}/characters/generate", json={"brief": "主角团三人", "count": 3}, ) assert gen_resp.status_code == 200 cards = gen_resp.json()["cards"] assert [c["name"] for c in cards] == ["叶无尘", "苏清歌", "胖三"] # primary 确实被调过且失败(耗尽重试 max_retries=1 → 2 次尝试)。 assert primary.complete_calls >= 1 # DB 真源断言:usage_ledger 记**实际服务方** = fallback(openai),不是 primary(deepseek)。 # served_by 不出 API,这是「回退真发生」的端到端 DB 证据(ARCH §4.5 记账记实际服务方)。 assert project_uuid is not None async with e2e_sm() as verify: ledger_rows = ( ( await verify.execute( select(UsageLedger).where(UsageLedger.project_id == project_uuid) ) ) .scalars() .all() ) assert len(ledger_rows) == 1 assert ledger_rows[0].provider == _FALLBACK assert ledger_rows[0].provider != _PRIMARY assert ledger_rows[0].input_tokens == _USAGE["writer"].input_tokens finally: if project_uuid is not None: await _cleanup(e2e_sm, project_uuid) async def test_m5_capability_degradation_marks_served_by_degraded() -> None: """用例 4:primary 不支持结构化输出且无支持者 → 网关降级用它,标 served_by.degraded(§4.4)。 聚焦集成测(直测真 `Gateway.run`,不经 HTTP——降级路径更清晰):结构化请求 (output_schema=CharacterGenResult)打到 `structured_output=False` 的唯一 provider → 降级(适配器自走 JSON 提示返 parsed),`served_by.degraded=True`。 """ from ww_llm_gateway.ledger import LedgerSink from ww_llm_gateway.types import Usage class _CapturingLedger: def __init__(self) -> None: self.records: list[Usage] = [] async def record(self, scope: Scope, usage: Usage) -> None: self.records.append(usage) # primary 不支持结构化输出,但仍能返回 parsed(自走 instructor JSON 提示)。 degrade_adapter = _FakeGenAdapter(provider=_PRIMARY, structured=False) ledger: LedgerSink = _CapturingLedger() # type: ignore[assignment] gw = Gateway( adapters={_PRIMARY: degrade_adapter}, ledger=ledger, chain_resolver=lambda _tier: [Route(provider=_PRIMARY, model="deepseek-chat")], ) req = LlmRequest( tier="writer", input="生成角色", output_schema=CharacterGenResult, scope=Scope(user_id=uuid.uuid4(), project_id=uuid.uuid4()), ) resp = await gw.run(req) assert isinstance(resp.parsed, CharacterGenResult) assert resp.served_by.provider == _PRIMARY # 链上无结构化支持者 → 降级用唯一可用 provider,标 degraded。 assert resp.served_by.degraded is True assert resp.served_by.fell_back is False