Files
writer-work-flow/tests/test_m5_e2e.py
Yaojia Wang 765dbdfbd4 feat: M4 文风 + M5 生成/多provider/Skill + Kimi Code 订阅接入 + 本地联调修复
M4(文风): style-auditor 双轨(提取指纹/漂移第四审)+ jobs 长任务框架(zombie reaper) + 回炉 refine + GET /style read-back。
M5(生成+扩展): worldbuilder/character-gen(入库 continuity 409 gate + partition_writes 白名单 + schema→JSONB 形变);
  网关多 provider 回退链/熔断/能力降级(Anthropic/Gemini 适配器);Skill registry + 表权限沙箱 + 规则;
  前端 角色生成器/世界观/Codex/规则页/技能库/⌘K 命令面板。
K1(Kimi Code 订阅接入): OAuth device-flow(kimi-code)+ 静态 Console key(kimi-code-key)两路径;
  coding 端点 KimiCLI 伪造头(实测 UA allow-list 门禁,缺则 403)+ JSON 模式结构化(thinking ⊥ tool_choice)。
本地联调修复: CORS 中间件;assemble 注入 premise+「写第N章」指令(修空 prompt 400);
  GET /outline·/draft read-back + 大纲/工作台/审稿页重载;写页 client/server 常量边界 + notFound 健壮化;
  字数 toLocaleString locale 水合;审稿页终稿从已存草稿 seed(修 accept 422)。
门禁: backend ruff/mypy(157)/alembic 无漂移/pytest 451 · frontend lint/tsc/vitest/build。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-20 10:39:58 +02:00

515 lines
22 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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) precheckanalyst 假网关产带冲突 `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`**
断言实际服务方是 **openaifallback而非 deepseekprimary**——证明回退真发生、记账记实际
服务方ARCH §4.5)。
用例 4 · 能力降级primary 不支持结构化输出 → 链上有支持者则重排/降级§4.4
primary(deepseek) `capabilities().structured_output=False` 且无 fallback → 网关降级用它
(适配器自走 JSON 提示返 parsed`served_by.degraded=True`。经真 `Gateway.run` 直测
(聚焦集成测,证明降级路径真走,不经 HTTP 更清晰)。
确定性 & 零成本(同 M1M4真实 `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` 选 schemaCharacterGenResult/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 默认都路由到 deepseekconfig.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 → 确认放行 → 真落 pgschema→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 == "灭门遗孤,背负血仇。"
# 记账真落 pgworld 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:
"""用例 3HTTP 生成 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:
"""用例 4primary 不支持结构化输出且无支持者 → 网关降级用它,标 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