"""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 · 三审 SSE(continuity/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` 路由 parsed(OutlineResult / 三审 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)默认都路由到 deepseek(config.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()`:吐固定写章 token(writer/draft),末尾喂 writer 用量。 - `complete()`:据 `req.output_schema` 路由: `OutlineResult`→大纲(analyst);`ContinuityReview`/`ForeshadowReview`/`PaceReview` →三审(analyst/light);`ChapterDigestFacts`→验收 digest(light,追溯终稿)。 每路带 `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:伏笔 埋设→进展→验收后到期扫描置 OVERDUE(BackgroundTask)。""" 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) 埋伏笔 F1(expected_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" # 对照:F2(expected_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 含 F1(status==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) → OVERDUE;progress 含那条进展。 assert by_code["F1"].status == "OVERDUE" assert progress_entry in by_code["F1"].progress # F2:expected_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:三审 SSE(continuity/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)