Files
writer-work-flow/tests/test_m3_e2e.py
Yaojia Wang 5fb7bfb1de feat: M3 — 伏笔账本 + 节奏引擎 + 大纲(含并发记账 bugfix)
- 伏笔账本:纯函数状态机(OPEN/PARTIAL/CLOSED/OVERDUE) + ForeshadowLedger repo;验收后到期扫描(BackgroundTask 自建 session 置 OVERDUE);登记/状态变更端点
- 节奏 + 三审齐:foreshadow-analyst + pace-checker 并入 LangGraph 并行审(REVIEW_SPECS),collect 分列落 chapter_reviews(conflicts/foreshadow_sug/pace),review SSE 加 foreshadow/pace 事件
- 大纲:outliner Agent 产 OutlineResult(含 foreshadow_windows),POST /outline 逐章 upsert outline 表;GET /foreshadow?status= 看板
- 前端:伏笔四泳道看板(OVERDUE 琥珀) + 大纲编辑器(窗口徽标) + 节奏节拍图(▁▃▅) + 审稿页消费 foreshadow/pace SSE
- bugfix(T3.8):并行三审共用请求 session 记账触发 'Session is already flushing' → foreshadow/pace 静默丢失;SqlAlchemyLedgerSink.record 改 add-only(靠端点/事务 commit),加并发回归测试
- M3 E2E:真实 DB + mock 网关零 token 走通 埋设→进展→验收后扫描 OVERDUE→看板 + 大纲含窗口 + 三审齐 SSE/留痕;E2E 暴露并钉住上述 bug
- 门禁绿:mypy 111 / pytest 228(0 xfailed) / alembic 无漂移;前端 gen:api/lint/tsc/vitest 69/build
2026-06-18 14:21:17 +02:00

585 lines
25 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.

"""M3 端到端:伏笔账本/到期提醒 + 大纲(含回收窗口) + 三审(真实 DB零 token
证明 M3 闭环(不变量 #2/#3/#4/#5
用例 1 · 伏笔 埋设→进展→逾期提醒(验收后扫描)
`POST /projects` → `POST .../foreshadow`(埋 F1, expected_close_to=2)
→ `PATCH .../foreshadow/F1`(追加一条进展) → `PUT .../chapters/5/draft`
→ `POST .../chapters/5/accept`(零冲突直通) → **验收后 BackgroundTask 到期扫描**置
`current_ch(5) > expected_close_to(2) AND status≠CLOSED → OVERDUE` →
`GET .../foreshadow?status=OVERDUE` 含 F1 + DB 真源断言F1.status==OVERDUE、progress 含那条)。
对照:未越界伏笔 F2(expected_close_to=9) 不应 OVERDUE。
用例 2 · 排大纲含伏笔回收窗口
`POST .../outline` → analyst 假适配器产 `OutlineResult` → DB `outline` 表逐章行存在
`beats` 形 `{"beats":[...]}`、`foreshadow_windows` JSONB list 含窗口);**再次 POST 幂等
覆盖**(唯一 `(project_id,chapter_no)` 行数不变);`usage_ledger` 含 analyst 那次记账。
用例 3 · 三审 SSEcontinuity/foreshadow/pace 三审齐 + 留痕列)
`POST .../chapters/N/review` → 解析帧含三审 section + ≥1 `foreshadow`/`pace` 事件;
`GET .../reviews` 留痕行 `foreshadow_sug`(list)/`pace`(dict) 列已填。
确定性 & 零成本(同 M1/M2真实 `Gateway` + 多档位假适配器(据 `req.output_schema`
分支返回固定 `parsed`,绝不联网)+ 真实 `SqlAlchemyLedgerSink`(用量记账闭环走通)。
无 pg → skip。
坑(见 memory/gotchas
- `LifespanManager` 触发 lifespan → `seed_stub_user`owner_id FK 依赖它)。
- `get_sessionmaker` engine 绑定首个 loop → 每 DB 测试清缓存重建、结束 dispose。
- 网关 ledger / 验收-side repos / outline upsert 均只 flush端点/验收事务 commit 才落库。
- **BackgroundTask 自建独立 session**FastAPI BackgroundTasks 在 response 发回、请求
session 关闭后才跑——`get_session_factory` 必被 override 成 `e2e_sm`(真 sessionmaker
否则 dep 解析会另建真 engine跨 loop 风险。httpx `ASGITransport` 下,`await client.post`
调用会等 ASGI app 协程(含 background tasks跑完才返回——故 client 上下文退出后断言稳定。
- 多档位假适配器据 `req.output_schema` 路由 parsedOutlineResult / 三审 schema /
ChapterDigestFacts写章走 stream。
"""
from __future__ import annotations
import json
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 (
ContinuityReview,
ForeshadowReview,
ForeshadowSuggestion,
ForeshadowWindow,
OutlineChapter,
OutlineResult,
PaceIssue,
PaceReview,
)
from ww_api.services.digest_extraction import ChapterDigestFacts
from ww_db import get_session, get_sessionmaker
from ww_db.models import (
Chapter,
ChapterDigest,
ChapterReview,
Foreshadow,
Outline,
Project,
UsageLedger,
)
from ww_llm_gateway import (
Gateway,
SqlAlchemyLedgerSink,
resolve_route,
)
from ww_llm_gateway.adapters.base import (
Capabilities,
ProviderResult,
ProviderUsage,
StreamChunk,
)
from ww_llm_gateway.types import LlmRequest
# 三档位writer/analyst/light默认都路由到 deepseekconfig.tier_defaults
_PROVIDER = "deepseek"
# 确定性流式写章 token。
_TOKENS = ["", "", "", "", "回收之时。"]
# 终稿里可被 digest 追溯的标记串(验证 digest 从终稿提炼,不变量 #4
_FINAL_MARKER = "终稿独有事件:旧伏笔在第五章浮现"
# outliner 假产物:两章,章 1 关联伏笔窗口(回收窗口 code=F1
_OUTLINE_RESULT = OutlineResult(
chapters=[
OutlineChapter(
no=1,
beats=["开端:少年入门", "埋线:神秘玉佩"],
foreshadow_windows=[
ForeshadowWindow(
code="F1",
plant_chapter=1,
expected_close_from=2,
expected_close_to=2,
)
],
),
OutlineChapter(
no=2,
beats=["试炼:初次交锋"],
foreshadow_windows=[],
),
]
)
# 三审固定产物。
_FORESHADOW_REVIEW = ForeshadowReview(
planted=[ForeshadowSuggestion(code="F9", title="新埋:山门禁制", where="第2段")],
resolved=[ForeshadowSuggestion(code="F1", title="疑似回收:玉佩", where="第5段")],
)
_PACE_REVIEW = PaceReview(
water=[PaceIssue(where="第3段", reason="信息密度低")],
hook=True,
beat_map=[1, 3, 5, 2],
)
# 各档位假用量(喂记账;证明各档位调用各自落库)。
_USAGE = {
"writer": ProviderUsage(input_tokens=11, output_tokens=5),
"analyst": ProviderUsage(input_tokens=23, output_tokens=7),
"light": ProviderUsage(input_tokens=17, output_tokens=9),
}
class _FakeMultiTierAdapter:
"""实现 `ProviderAdapter` Protocol按 `req.output_schema` 分支返回 parsed绝不联网。
- `stream()`:吐固定写章 tokenwriter/draft末尾喂 writer 用量。
- `complete()`:据 `req.output_schema` 路由:
`OutlineResult`→大纲analyst`ContinuityReview`/`ForeshadowReview`/`PaceReview`
→三审analyst/light`ChapterDigestFacts`→验收 digestlight追溯终稿
每路带 `ProviderUsage` → 网关据此各落一条 usage_ledger。
"""
provider = _PROVIDER
def capabilities(self) -> Capabilities:
return Capabilities(structured_output=True)
async def complete(self, req: LlmRequest, model: str) -> ProviderResult:
schema = req.output_schema
if schema is OutlineResult:
return ProviderResult(
text=_OUTLINE_RESULT.model_dump_json(),
parsed=_OUTLINE_RESULT,
usage=_USAGE["analyst"],
)
if schema is ContinuityReview:
review = ContinuityReview(conflicts=[]) # 零冲突 → 验收 gate 直通
return ProviderResult(
text=review.model_dump_json(),
parsed=review,
usage=_USAGE["analyst"],
)
if schema is ForeshadowReview:
return ProviderResult(
text=_FORESHADOW_REVIEW.model_dump_json(),
parsed=_FORESHADOW_REVIEW,
usage=_USAGE["analyst"],
)
if schema is PaceReview:
return ProviderResult(
text=_PACE_REVIEW.model_dump_json(),
parsed=_PACE_REVIEW,
usage=_USAGE["light"],
)
# 否则视为 digest 提炼ChapterDigestFacts把终稿标记塞进 summary#4
facts = ChapterDigestFacts(
summary=_FINAL_MARKER,
events=["玉佩浮现"],
locations=["山门"],
)
return ProviderResult(
text=facts.model_dump_json(),
parsed=facts,
usage=_USAGE["light"],
)
async def stream(self, req: LlmRequest, model: str) -> AsyncIterator[StreamChunk]:
for token in _TOKENS:
yield StreamChunk(text=token)
yield StreamChunk(usage=_USAGE["writer"])
@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 _parse_sse(raw: str) -> list[tuple[str, str]]:
"""把 text/event-stream 原文解析为 `(event, data)` 帧列表。"""
frames: list[tuple[str, str]] = []
event: str | None = None
data: str | None = None
for line in raw.splitlines():
if line.startswith("event:"):
event = line[len("event:") :].strip()
elif line.startswith("data:"):
data = line[len("data:") :].strip()
elif line == "":
if event is not None and data is not None:
frames.append((event, data))
event, data = None, None
if event is not None and data is not None:
frames.append((event, data))
return frames
def _gateway_override_factory(
adapter: _FakeMultiTierAdapter,
) -> Callable[[AsyncSession], Gateway]:
"""构造依赖覆盖:真实 Gateway + 给定假适配器 + 真实 ledger请求 session
ledger 用**请求 session**`Depends(get_session)`FastAPI 按请求缓存,与端点/验收
事务同一实例)→ 由端点流末 / 验收事务的 `commit()` 落库(同 M1/M2 回归校验)。
"""
def _override(session: Annotated[AsyncSession, Depends(get_session)]) -> Gateway:
return Gateway(
adapters={_PROVIDER: adapter},
ledger=SqlAlchemyLedgerSink(session),
resolver=resolve_route,
)
return _override
async def _cleanup(e2e_sm: async_sessionmaker[AsyncSession], project_uuid: uuid.UUID) -> None:
"""按 FK 顺序清理(无级联的先删,再删 project多数子表经 FK CASCADE"""
async with e2e_sm() as cleanup:
await cleanup.execute(delete(UsageLedger).where(UsageLedger.project_id == project_uuid))
await cleanup.execute(delete(ChapterDigest).where(ChapterDigest.project_id == project_uuid))
await cleanup.execute(delete(ChapterReview).where(ChapterReview.project_id == project_uuid))
await cleanup.execute(delete(Outline).where(Outline.project_id == project_uuid))
await cleanup.execute(delete(Foreshadow).where(Foreshadow.project_id == project_uuid))
await cleanup.execute(delete(Chapter).where(Chapter.project_id == project_uuid))
await cleanup.execute(delete(Project).where(Project.id == project_uuid))
await cleanup.commit()
async def test_m3_foreshadow_plant_progress_overdue_after_accept(
e2e_sm: async_sessionmaker[AsyncSession],
) -> None:
"""用例 1伏笔 埋设→进展→验收后到期扫描置 OVERDUEBackgroundTask"""
from ww_api.main import create_app
from ww_api.services.project_deps import (
get_digest_gateway,
get_review_gateway,
get_session_factory,
get_writer_gateway,
)
adapter = _FakeMultiTierAdapter()
override = _gateway_override_factory(adapter)
app = create_app()
app.dependency_overrides[get_writer_gateway] = override
app.dependency_overrides[get_review_gateway] = override
app.dependency_overrides[get_digest_gateway] = override
# 到期扫描 BackgroundTask 的独立 session 工厂 → 注入真实 sessionmaker同测试 engine/loop
# 否则默认 `get_sessionmaker()` 会另建 engine 绑别的 loop跨 loop 坑)。
app.dependency_overrides[get_session_factory] = lambda: e2e_sm
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:
# 1) 立项 → 201。
create_resp = await client.post(
"/projects",
json={"title": "M3 伏笔逾期验证作品", "genre": "玄幻"},
)
assert create_resp.status_code == 201
project_id = create_resp.json()["id"]
project_uuid = uuid.UUID(project_id)
# 2) 埋伏笔 F1expected_close_to=2 → 验收第5章后应逾期→ 201, status=OPEN。
f1_resp = await client.post(
f"/projects/{project_id}/foreshadow",
json={
"code": "F1",
"title": "神秘玉佩",
"planted_at": 1,
"expected_close_to": 2,
},
)
assert f1_resp.status_code == 201
assert f1_resp.json()["status"] == "OPEN"
# 对照F2expected_close_to=9 → 不越界,不应逾期)。
f2_resp = await client.post(
f"/projects/{project_id}/foreshadow",
json={"code": "F2", "title": "远期线索", "expected_close_to": 9},
)
assert f2_resp.status_code == 201
# 3) PATCH F1 追加一条进展 → 200, progress 追加。
progress_entry = {"chapter": 3, "note": "玉佩发热"}
patch_resp = await client.patch(
f"/projects/{project_id}/foreshadow/F1",
json={"progress_entry": progress_entry},
)
assert patch_resp.status_code == 200
assert patch_resp.json()["progress"][-1] == progress_entry
# 4) 写第5章草稿SSE+ 自动保存。
draft_resp = await client.post(f"/projects/{project_id}/chapters/5/draft")
assert draft_resp.status_code == 200
draft_text = "".join(
json.loads(d)["text"]
for (ev, d) in _parse_sse(draft_resp.text)
if ev == "token"
)
assert draft_text == "".join(_TOKENS)
save_resp = await client.put(
f"/projects/{project_id}/chapters/5/draft",
json={"text": draft_text},
)
assert save_resp.status_code == 200
# 5) 验收第5章无审稿留痕 → 零冲突直通final_text 含终稿标记 #4
final_text = draft_text + "\n\n" + _FINAL_MARKER
accept_resp = await client.post(
f"/projects/{project_id}/chapters/5/accept",
json={"final_text": final_text, "decisions": []},
)
assert accept_resp.status_code == 200
acc = accept_resp.json()
assert acc["digest_added"] is True
# BackgroundTask 到期扫描在此请求 ASGI 协程内response 后)跑完才返回。
# 6a) 端点视角GET ?status=OVERDUE 含 F1status==OVERDUE不含 F2。
board_resp = await client.get(
f"/projects/{project_id}/foreshadow", params={"status": "OVERDUE"}
)
assert board_resp.status_code == 200
overdue = board_resp.json()["foreshadow"]
overdue_codes = {f["code"] for f in overdue}
assert "F1" in overdue_codes
assert "F2" not in overdue_codes
f1_board = next(f for f in overdue if f["code"] == "F1")
assert f1_board["status"] == "OVERDUE"
# 6b) DB 真源断言(经 e2e session
assert project_uuid is not None
async with e2e_sm() as verify:
fore_rows = (
(
await verify.execute(
select(Foreshadow)
.where(Foreshadow.project_id == project_uuid)
.order_by(Foreshadow.code)
)
)
.scalars()
.all()
)
by_code = {r.code: r for r in fore_rows}
# F1验收第5章 > expected_close_to(2) → OVERDUEprogress 含那条进展。
assert by_code["F1"].status == "OVERDUE"
assert progress_entry in by_code["F1"].progress
# F2expected_close_to(9) 未越界 → 仍 OPEN对照扫描不误伤
assert by_code["F2"].status == "OPEN"
# 终稿 digest 落库且 summary 含终稿标记(#4从终稿非草稿提炼
digest = (
await verify.execute(
select(ChapterDigest).where(ChapterDigest.project_id == project_uuid)
)
).scalar_one()
assert digest.facts["summary"] == _FINAL_MARKER
assert _FINAL_MARKER not in draft_text
finally:
if project_uuid is not None:
await _cleanup(e2e_sm, project_uuid)
async def test_m3_outline_with_foreshadow_windows_upsert_idempotent(
e2e_sm: async_sessionmaker[AsyncSession],
) -> None:
"""用例 2排大纲含伏笔回收窗口逐章落 `outline`;再次 POST 幂等覆盖;记账落库。"""
from ww_api.main import create_app
from ww_api.services.project_deps import get_outline_gateway
adapter = _FakeMultiTierAdapter()
override = _gateway_override_factory(adapter)
app = create_app()
app.dependency_overrides[get_outline_gateway] = override
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": "M3 大纲窗口验证作品"})
assert create_resp.status_code == 201
project_id = create_resp.json()["id"]
project_uuid = uuid.UUID(project_id)
# 1) 生成大纲analyst 假适配器产 OutlineResult→ 200 OutlineResponse。
outline_resp = await client.post(
f"/projects/{project_id}/outline", json={"volume": 1}
)
assert outline_resp.status_code == 200
chapters = outline_resp.json()["chapters"]
assert len(chapters) == 2
ch1 = next(c for c in chapters if c["no"] == 1)
assert ch1["volume"] == 1
assert ch1["beats"] == ["开端:少年入门", "埋线:神秘玉佩"]
assert ch1["foreshadow_windows"][0]["code"] == "F1"
assert ch1["foreshadow_windows"][0]["expected_close_to"] == 2
# 2) 再次 POST → 唯一 (project_id, chapter_no) 覆盖,行数不堆叠。
outline_resp2 = await client.post(
f"/projects/{project_id}/outline", json={"volume": 1}
)
assert outline_resp2.status_code == 200
# DB 真源断言。
assert project_uuid is not None
async with e2e_sm() as verify:
outline_rows = (
(
await verify.execute(
select(Outline)
.where(Outline.project_id == project_uuid)
.order_by(Outline.chapter_no)
)
)
.scalars()
.all()
)
# 幂等:两次 POST 仍是 2 行(每章一行,覆盖非堆叠)。
assert len(outline_rows) == 2
row1 = outline_rows[0]
assert row1.chapter_no == 1
assert row1.volume == 1
# beats 列形 `{"beats":[...]}`(写侧包裹,见 contracts C3 扩 T3.5)。
assert row1.beats == {"beats": ["开端:少年入门", "埋线:神秘玉佩"]}
# foreshadow_windows JSONB list 含窗口code=F1
assert isinstance(row1.foreshadow_windows, list)
assert row1.foreshadow_windows[0]["code"] == "F1"
assert row1.foreshadow_windows[0]["expected_close_to"] == 2
# 章 2 无窗口。
assert outline_rows[1].chapter_no == 2
assert outline_rows[1].foreshadow_windows == []
# usage_ledger两次大纲生成各 1 条 analyst 调用记账(端点 commit 后落库)。
ledger_rows = (
(
await verify.execute(
select(UsageLedger).where(UsageLedger.project_id == project_uuid)
)
)
.scalars()
.all()
)
assert len(ledger_rows) == 2
assert all(r.input_tokens == _USAGE["analyst"].input_tokens for r in ledger_rows)
assert all(r.provider == _PROVIDER for r in ledger_rows)
finally:
if project_uuid is not None:
await _cleanup(e2e_sm, project_uuid)
async def test_m3_three_reviewers_sse_and_review_trace(
e2e_sm: async_sessionmaker[AsyncSession],
) -> None:
"""用例 3三审 SSEcontinuity/foreshadow/pace 三审齐)+ 留痕列foreshadow_sug/pace真落库。
三审在同一 LangGraph superstep 并行跑、共用请求 session 记账T3.8 把
`SqlAlchemyLedgerSink.record` 改 add-only`session.add` 同步不让步 → 并行协程不交错,
持久化靠端点 commit。故三审齐现真持久化断言严格正向
- SSE `section` 帧含 continuity/foreshadow/pace 三审,且无 `error` 帧;
- SSE 至少各一条 `foreshadow` / `pace` 事件foreshadow/pace 真发事件);
- `GET .../reviews` 留痕行 `foreshadow_sug`(list 非空) / `pace`(dict 非空) 列已填;
- DB 真源 `chapter_reviews.foreshadow_sug`/`pace` 列真落库(三审都 done无静默丢失
"""
from ww_api.main import create_app
from ww_api.services.project_deps import (
get_review_gateway,
get_writer_gateway,
)
adapter = _FakeMultiTierAdapter()
override = _gateway_override_factory(adapter)
app = create_app()
app.dependency_overrides[get_writer_gateway] = override
app.dependency_overrides[get_review_gateway] = override
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": "M3 三审验证作品"})
assert create_resp.status_code == 201
project_id = create_resp.json()["id"]
project_uuid = uuid.UUID(project_id)
# 写第1章草稿 + 自动保存(供续审回退已存草稿)。
draft_resp = await client.post(f"/projects/{project_id}/chapters/1/draft")
draft_text = "".join(
json.loads(d)["text"]
for (ev, d) in _parse_sse(draft_resp.text)
if ev == "token"
)
await client.put(
f"/projects/{project_id}/chapters/1/draft",
json={"text": draft_text},
)
# 续审SSE三审 section + ≥1 foreshadow/pace 事件、无 error。
review_resp = await client.post(
f"/projects/{project_id}/chapters/1/review",
json={"draft": draft_text},
)
assert review_resp.status_code == 200
assert review_resp.headers["content-type"].startswith("text/event-stream")
rframes = _parse_sse(review_resp.text)
assert [d for (ev, d) in rframes if ev == "error"] == []
section_names = {
json.loads(d).get("name") for (ev, d) in rframes if ev == "section"
}
assert {"continuity", "foreshadow", "pace"} <= section_names
foreshadow_events = [d for (ev, d) in rframes if ev == "foreshadow"]
pace_events = [d for (ev, d) in rframes if ev == "pace"]
assert len(foreshadow_events) >= 1
assert len(pace_events) >= 1
# GET 历史:留痕行 foreshadow_sug(list) / pace(dict) 列已填。
hist_resp = await client.get(f"/projects/{project_id}/chapters/1/reviews")
assert hist_resp.status_code == 200
latest = hist_resp.json()["reviews"][0]
assert isinstance(latest["foreshadow_sug"], list)
assert len(latest["foreshadow_sug"]) >= 1
assert isinstance(latest["pace"], dict)
assert latest["pace"]["hook"] is True
# DB 真源断言chapter_reviews 留痕列已填。
assert project_uuid is not None
async with e2e_sm() as verify:
review_rows = (
(
await verify.execute(
select(ChapterReview).where(ChapterReview.project_id == project_uuid)
)
)
.scalars()
.all()
)
assert len(review_rows) == 1
row = review_rows[0]
assert isinstance(row.foreshadow_sug, list)
assert len(row.foreshadow_sug) >= 1
# 扁平后每条带 kind:"planted"|"resolved"(见 collect 列映射)。
assert {s.get("kind") for s in row.foreshadow_sug} <= {"planted", "resolved"}
assert isinstance(row.pace, dict)
assert row.pace["hook"] is True
assert row.pace["beat_map"] == [1, 3, 5, 2]
finally:
if project_uuid is not None:
await _cleanup(e2e_sm, project_uuid)