assemble/_build_volatile 加 directive 参数(仅入 volatile,守缓存前缀不变量 #9);stream_draft 加可选 body {directive} 透传;不持久化、无迁移;directive 长度入 log 不记原文。新 DraftStreamRequest schema + 重生成 TS 客户端。后端门禁绿: ruff/mypy/pytest。契约记 memory/contracts.md。
530 lines
20 KiB
Python
530 lines
20 KiB
Python
"""项目(立项)+ 章节草稿端点(C3 / ARCH §7.2, §7.3;不变量 #7)。
|
||
|
||
- POST /projects 立项向导 → projects 行(owner_id=stub)。
|
||
- GET /projects 列出项目。
|
||
- GET /projects/:id 项目详情(404 → NOT_FOUND)。
|
||
- POST /projects/:id/chapters/:no/draft 流式写章草稿(SSE,text/event-stream)。
|
||
- PUT /projects/:id/chapters/:no/draft 自动保存草稿(幂等 upsert)。
|
||
|
||
不变量:writer 只产草稿(不提升 accepted 版本/不抽 digest/不写状态——属 M2 验收);
|
||
agent 只传 tier;DB 按 project_id/owner_id 过滤;日志脱敏(只记长度,不记正文)。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import uuid
|
||
from collections.abc import AsyncIterator
|
||
from typing import Annotated
|
||
|
||
from fastapi import APIRouter, BackgroundTasks, Depends, Request
|
||
from fastapi.responses import StreamingResponse
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
from ww_core.domain.chapter_repo import ChapterRepo
|
||
from ww_core.domain.digest_repo import DigestAppendRepo
|
||
from ww_core.domain.injection_repo import EntityRef, InjectionOverride, InjectionOverrideRepo
|
||
from ww_core.domain.project_repo import ProjectCreate, ProjectRepo
|
||
from ww_core.domain.repositories import MemoryRepos
|
||
from ww_core.domain.review_repo import ReviewRepo
|
||
from ww_core.memory import assemble
|
||
from ww_core.memory.selection import RECENT_DIGEST_COUNT
|
||
from ww_core.orchestrator import (
|
||
ChapterState,
|
||
SseEvent,
|
||
build_review_context,
|
||
build_review_graph,
|
||
normalize_deltas,
|
||
normalize_review,
|
||
stream_chapter_draft,
|
||
)
|
||
from ww_db import get_session
|
||
from ww_llm_gateway import Gateway
|
||
from ww_shared import AppError, ErrorCode
|
||
|
||
from ww_api.logging_config import get_logger
|
||
from ww_api.schemas.injection import (
|
||
InjectionEntity,
|
||
InjectionEntityRef,
|
||
InjectionOverrideRequest,
|
||
InjectionResponse,
|
||
)
|
||
from ww_api.schemas.projects import (
|
||
AcceptRequest,
|
||
AcceptResponse,
|
||
DraftResponse,
|
||
DraftSaveRequest,
|
||
DraftStreamRequest,
|
||
DraftView,
|
||
ProjectCreateRequest,
|
||
ProjectListResponse,
|
||
ProjectResponse,
|
||
ReviewHistoryItem,
|
||
ReviewHistoryResponse,
|
||
ReviewRequest,
|
||
)
|
||
from ww_api.services.accept_service import (
|
||
AcceptOutcome,
|
||
assert_conflicts_resolved,
|
||
run_accept_transaction,
|
||
)
|
||
from ww_api.services.credentials import STUB_OWNER_ID
|
||
from ww_api.services.digest_extraction import extract_digest_facts
|
||
from ww_api.services.foreshadow_scan import SessionFactory, run_overdue_scan
|
||
from ww_api.services.project_deps import (
|
||
get_chapter_repo,
|
||
get_digest_append_repo,
|
||
get_digest_gateway,
|
||
get_injection_repo,
|
||
get_memory_repos,
|
||
get_project_repo,
|
||
get_review_gateway,
|
||
get_review_repo,
|
||
get_session_factory,
|
||
get_writer_gateway,
|
||
)
|
||
|
||
log = get_logger("ww.api.projects")
|
||
|
||
router = APIRouter(prefix="/projects", tags=["projects"])
|
||
|
||
ProjectRepoDep = Annotated[ProjectRepo, Depends(get_project_repo)]
|
||
ChapterRepoDep = Annotated[ChapterRepo, Depends(get_chapter_repo)]
|
||
GatewayDep = Annotated[Gateway, Depends(get_writer_gateway)]
|
||
ReviewGatewayDep = Annotated[Gateway, Depends(get_review_gateway)]
|
||
DigestGatewayDep = Annotated[Gateway, Depends(get_digest_gateway)]
|
||
MemoryReposDep = Annotated[MemoryRepos, Depends(get_memory_repos)]
|
||
InjectionRepoDep = Annotated[InjectionOverrideRepo, Depends(get_injection_repo)]
|
||
ReviewRepoDep = Annotated[ReviewRepo, Depends(get_review_repo)]
|
||
DigestRepoDep = Annotated[DigestAppendRepo, Depends(get_digest_append_repo)]
|
||
|
||
|
||
def _to_response(view: object) -> ProjectResponse:
|
||
# ProjectView 与 ProjectResponse 字段同名,逐字段映射(snake_case 契约)。
|
||
return ProjectResponse.model_validate(view, from_attributes=True)
|
||
|
||
|
||
@router.post("", status_code=201)
|
||
async def create_project(body: ProjectCreateRequest, repo: ProjectRepoDep) -> ProjectResponse:
|
||
view = await repo.create(
|
||
STUB_OWNER_ID,
|
||
ProjectCreate(
|
||
title=body.title,
|
||
genre=body.genre,
|
||
logline=body.logline,
|
||
premise=body.premise,
|
||
theme=body.theme,
|
||
selling_points=body.selling_points,
|
||
structure=body.structure,
|
||
),
|
||
)
|
||
log.info("project_created", project_id=str(view.id), title_len=len(view.title))
|
||
return _to_response(view)
|
||
|
||
|
||
@router.get("")
|
||
async def list_projects(repo: ProjectRepoDep) -> ProjectListResponse:
|
||
views = await repo.list_for_owner(STUB_OWNER_ID)
|
||
return ProjectListResponse(projects=[_to_response(v) for v in views])
|
||
|
||
|
||
@router.get("/{project_id}")
|
||
async def get_project(project_id: uuid.UUID, repo: ProjectRepoDep) -> ProjectResponse:
|
||
view = await repo.get(STUB_OWNER_ID, project_id)
|
||
if view is None:
|
||
raise AppError(ErrorCode.NOT_FOUND, f"project {project_id} not found")
|
||
return _to_response(view)
|
||
|
||
|
||
async def _injection_response(
|
||
repos: MemoryRepos,
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
override: InjectionOverride | None,
|
||
) -> InjectionResponse:
|
||
"""组装注入响应:以作者覆盖跑确定性 assemble → selected + override 回显(GET/PUT 共用)。"""
|
||
context = await assemble(repos, project_id, chapter_no, override=override)
|
||
selected = [
|
||
InjectionEntity(kind=e.kind, name=e.name, reasons=list(e.reasons))
|
||
for e in context.selection.selected
|
||
]
|
||
effective_n = override.recent_n if (override and override.recent_n) else RECENT_DIGEST_COUNT
|
||
return InjectionResponse(
|
||
project_id=project_id,
|
||
chapter_no=chapter_no,
|
||
selected=selected,
|
||
recent_n=effective_n,
|
||
pinned=[
|
||
InjectionEntityRef(kind=r.kind, name=r.name)
|
||
for r in (override.pinned if override else [])
|
||
],
|
||
excluded=[
|
||
InjectionEntityRef(kind=r.kind, name=r.name)
|
||
for r in (override.excluded if override else [])
|
||
],
|
||
)
|
||
|
||
|
||
@router.get("/{project_id}/chapters/{chapter_no}/injection")
|
||
async def get_injection(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
repos: MemoryReposDep,
|
||
project_repo: ProjectRepoDep,
|
||
injection_repo: InjectionRepoDep,
|
||
) -> InjectionResponse:
|
||
"""本章注入透明(B0,读端点):回放确定性 SelectionTrace + 作者覆盖。无 LLM、无 commit。
|
||
|
||
项目不存在 → 404;无大纲 → selected: [](不报错)。看到的=写章用的(不变量 #6)。
|
||
"""
|
||
if await project_repo.get(STUB_OWNER_ID, project_id) is None:
|
||
raise AppError(ErrorCode.NOT_FOUND, f"project {project_id} not found")
|
||
override = await injection_repo.get(project_id, chapter_no)
|
||
resp = await _injection_response(repos, project_id, chapter_no, override)
|
||
log.info(
|
||
"injection_read",
|
||
project_id=str(project_id),
|
||
chapter_no=chapter_no,
|
||
selected_count=len(resp.selected),
|
||
pinned_count=len(resp.pinned),
|
||
excluded_count=len(resp.excluded),
|
||
)
|
||
return resp
|
||
|
||
|
||
@router.put("/{project_id}/chapters/{chapter_no}/injection")
|
||
async def save_injection(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
body: InjectionOverrideRequest,
|
||
repos: MemoryReposDep,
|
||
project_repo: ProjectRepoDep,
|
||
injection_repo: InjectionRepoDep,
|
||
session: Annotated[AsyncSession, Depends(get_session)],
|
||
) -> InjectionResponse:
|
||
"""本章注入覆盖(B0 可控版,写端点):存 pin/排除/recent_n → 回放新的确定性 selected。
|
||
|
||
项目不存在 → 404。upsert 只 flush,端点末尾 commit。覆盖是确定性输入,draft 端点同读,
|
||
故「看到的=写章用的」仍成立(不变量 #6 的作者兜底)。
|
||
"""
|
||
if await project_repo.get(STUB_OWNER_ID, project_id) is None:
|
||
raise AppError(ErrorCode.NOT_FOUND, f"project {project_id} not found")
|
||
override = InjectionOverride(
|
||
pinned=[EntityRef(kind=r.kind, name=r.name) for r in body.pinned],
|
||
excluded=[EntityRef(kind=r.kind, name=r.name) for r in body.excluded],
|
||
recent_n=body.recent_n,
|
||
)
|
||
saved = await injection_repo.upsert(project_id, chapter_no, override)
|
||
await session.commit()
|
||
log.info(
|
||
"injection_saved",
|
||
project_id=str(project_id),
|
||
chapter_no=chapter_no,
|
||
pinned_count=len(saved.pinned),
|
||
excluded_count=len(saved.excluded),
|
||
recent_n=saved.recent_n,
|
||
)
|
||
return await _injection_response(repos, project_id, chapter_no, saved)
|
||
|
||
|
||
def _encode_sse(event: SseEvent) -> str:
|
||
"""把归一事件编码为 text/event-stream 帧:`event: <name>\\ndata: <json>\\n\\n`。"""
|
||
payload = json.dumps(event.data, ensure_ascii=False)
|
||
return f"event: {event.event}\ndata: {payload}\n\n"
|
||
|
||
|
||
@router.post("/{project_id}/chapters/{chapter_no}/draft")
|
||
async def stream_draft(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
request: Request,
|
||
repos: MemoryReposDep,
|
||
gateway: GatewayDep,
|
||
injection_repo: InjectionRepoDep,
|
||
session: Annotated[AsyncSession, Depends(get_session)],
|
||
body: DraftStreamRequest | None = None,
|
||
) -> StreamingResponse:
|
||
"""流式写章草稿:组装记忆(含作者注入覆盖 + 本章指令)→ 网关流 → 归一 SSE → event-stream。
|
||
|
||
`body.directive`(可选,T4-b)是临时本章指令,直通 assemble→volatile(不持久化);
|
||
无 body 的旧调用方仍可用(向后兼容)。
|
||
"""
|
||
request_id = getattr(request.state, "request_id", None)
|
||
directive = body.directive if body else None
|
||
override = await injection_repo.get(project_id, chapter_no)
|
||
context = await assemble(repos, project_id, chapter_no, override=override, directive=directive)
|
||
log.info(
|
||
"draft_stream_start",
|
||
project_id=str(project_id),
|
||
chapter_no=chapter_no,
|
||
request_id=request_id,
|
||
stable_core_len=len(context.stable_core),
|
||
volatile_len=len(context.volatile),
|
||
directive_len=len(body.directive) if body and body.directive else 0,
|
||
)
|
||
|
||
deltas = stream_chapter_draft(
|
||
gateway,
|
||
stable_core=context.stable_core,
|
||
volatile=context.volatile,
|
||
user_id=STUB_OWNER_ID,
|
||
project_id=project_id,
|
||
)
|
||
|
||
async def _frames() -> AsyncIterator[str]:
|
||
async for event in normalize_deltas(deltas, request_id=request_id):
|
||
yield _encode_sse(event)
|
||
# 网关在流末经 SqlAlchemyLedgerSink.record 把 usage_ledger 行 flush 进本请求 session;
|
||
# sink 按设计不提交(写库事务由编排层控制,见不变量)。draft 端点无其他写副作用,
|
||
# 故流耗尽后在此提交,确保「每次调用一条 usage_ledger」真正落库(T1.9 暴露)。
|
||
await session.commit()
|
||
|
||
return StreamingResponse(
|
||
_frames(),
|
||
media_type="text/event-stream",
|
||
headers={"cache-control": "no-cache", "x-accel-buffering": "no"},
|
||
)
|
||
|
||
|
||
@router.put("/{project_id}/chapters/{chapter_no}/draft")
|
||
async def save_draft(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
body: DraftSaveRequest,
|
||
repo: ChapterRepoDep,
|
||
) -> DraftResponse:
|
||
"""自动保存:幂等 upsert 草稿(同章节覆盖同一行,版次不爆炸)。"""
|
||
view = await repo.save_draft(project_id, chapter_no, text=body.text)
|
||
log.info(
|
||
"draft_saved",
|
||
project_id=str(project_id),
|
||
chapter_no=chapter_no,
|
||
length=len(view.content),
|
||
)
|
||
return DraftResponse(
|
||
project_id=view.project_id,
|
||
chapter_no=view.chapter_no,
|
||
volume=view.volume,
|
||
status=view.status,
|
||
version=view.version,
|
||
length=len(view.content),
|
||
)
|
||
|
||
|
||
@router.get("/{project_id}/chapters/{chapter_no}/draft")
|
||
async def get_draft(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
repo: ChapterRepoDep,
|
||
) -> DraftView:
|
||
"""读取已保存草稿(含正文),供工作台重访时重载编辑器(镜像 GET /outline 读侧)。
|
||
|
||
复用 `chapter_repo.get_draft`(与续审 `_resolve_review_draft` 同一读 seam);只读不写库。
|
||
无草稿行(含空正文)→ 404 NOT_FOUND,工作台据此呈现空编辑器(与其它「缺资源」端点一致)。
|
||
多版本时该读取固定返草稿版次(DRAFT_VERSION)这条可编辑工作副本,与 PUT 保存的同一行。
|
||
"""
|
||
view = await repo.get_draft(project_id, chapter_no)
|
||
if view is None or not view.content.strip():
|
||
raise AppError(
|
||
ErrorCode.NOT_FOUND,
|
||
f"chapter {chapter_no} has no saved draft",
|
||
)
|
||
log.info(
|
||
"draft_read",
|
||
project_id=str(project_id),
|
||
chapter_no=chapter_no,
|
||
length=len(view.content),
|
||
)
|
||
return DraftView(
|
||
project_id=view.project_id,
|
||
chapter_no=view.chapter_no,
|
||
volume=view.volume,
|
||
status=view.status,
|
||
version=view.version,
|
||
content=view.content,
|
||
length=len(view.content),
|
||
)
|
||
|
||
|
||
async def _resolve_review_draft(
|
||
body: ReviewRequest,
|
||
chapter_repo: ChapterRepo,
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
) -> str:
|
||
"""取待审正文:请求体 `draft` 优先;否则回退到已保存草稿;都无 → NOT_FOUND。"""
|
||
if body.draft is not None and body.draft.strip():
|
||
return body.draft
|
||
saved = await chapter_repo.get_draft(project_id, chapter_no)
|
||
if saved is None or not saved.content.strip():
|
||
raise AppError(
|
||
ErrorCode.NOT_FOUND,
|
||
f"chapter {chapter_no} has no draft to review; provide `draft` in body",
|
||
)
|
||
return saved.content
|
||
|
||
|
||
@router.post("/{project_id}/chapters/{chapter_no}/review")
|
||
async def review_chapter(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
body: ReviewRequest,
|
||
request: Request,
|
||
repos: MemoryReposDep,
|
||
chapter_repo: ChapterRepoDep,
|
||
review_repo: ReviewRepoDep,
|
||
gateway: ReviewGatewayDep,
|
||
session: Annotated[AsyncSession, Depends(get_session)],
|
||
) -> StreamingResponse:
|
||
"""续审(SSE):组审稿上下文 → 跑审稿子图 → 归一为 section/conflict/done 事件。
|
||
|
||
提交边界:网关 ledger + collect 经 review_repo.record 均只 flush;端点在**流耗尽后**
|
||
`await session.commit()`(镜像 draft 端点,否则记账/留痕静默丢失)。
|
||
"""
|
||
request_id = getattr(request.state, "request_id", None)
|
||
draft = await _resolve_review_draft(body, chapter_repo, project_id, chapter_no)
|
||
context = await assemble(repos, project_id, chapter_no)
|
||
review_context = build_review_context(
|
||
draft=draft, stable_core=context.stable_core, volatile=context.volatile
|
||
)
|
||
log.info(
|
||
"review_stream_start",
|
||
project_id=str(project_id),
|
||
chapter_no=chapter_no,
|
||
request_id=request_id,
|
||
draft_len=len(draft),
|
||
review_context_len=len(review_context),
|
||
)
|
||
|
||
graph = build_review_graph(gateway, review_repo)
|
||
initial: ChapterState = {
|
||
"project_id": project_id,
|
||
"chapter_no": chapter_no,
|
||
"user_id": STUB_OWNER_ID,
|
||
"review_context": review_context,
|
||
}
|
||
|
||
async def _frames() -> AsyncIterator[str]:
|
||
final = await graph.ainvoke(initial)
|
||
reviews = final.get("reviews") or {}
|
||
async for event in normalize_review(reviews, request_id=request_id):
|
||
yield _encode_sse(event)
|
||
# collect 经 review_repo.record 落 chapter_reviews(只 flush)+ 网关 ledger 只 flush
|
||
# → 流耗尽后在此提交,确保审稿留痕 + usage_ledger 真正落库(同 draft 端点)。
|
||
await session.commit()
|
||
|
||
return StreamingResponse(
|
||
_frames(),
|
||
media_type="text/event-stream",
|
||
headers={"cache-control": "no-cache", "x-accel-buffering": "no"},
|
||
)
|
||
|
||
|
||
@router.get("/{project_id}/chapters/{chapter_no}/reviews")
|
||
async def list_reviews(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
review_repo: ReviewRepoDep,
|
||
) -> ReviewHistoryResponse:
|
||
"""审稿历史(新→旧):供前端审稿页加载既往审稿留痕 + 裁决。"""
|
||
views = await review_repo.list_for_chapter(project_id, chapter_no)
|
||
items = [
|
||
ReviewHistoryItem(
|
||
id=v.id,
|
||
project_id=v.project_id,
|
||
chapter_no=v.chapter_no,
|
||
chapter_version=v.chapter_version,
|
||
conflicts=v.conflicts,
|
||
foreshadow_sug=v.foreshadow_sug,
|
||
style=v.style,
|
||
pace=v.pace,
|
||
health_score=v.health_score,
|
||
decisions=v.decisions,
|
||
)
|
||
for v in views
|
||
]
|
||
return ReviewHistoryResponse(reviews=items)
|
||
|
||
|
||
@router.post("/{project_id}/chapters/{chapter_no}/accept")
|
||
async def accept_chapter(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
body: AcceptRequest,
|
||
request: Request,
|
||
background_tasks: BackgroundTasks,
|
||
chapter_repo: ChapterRepoDep,
|
||
digest_repo: DigestRepoDep,
|
||
review_repo: ReviewRepoDep,
|
||
gateway: DigestGatewayDep,
|
||
session: Annotated[AsyncSession, Depends(get_session)],
|
||
session_factory: Annotated[SessionFactory, Depends(get_session_factory)],
|
||
) -> AcceptResponse:
|
||
"""验收事务 + 冲突 gate(§5.5):gate(事务前)→ 终稿提炼 digest(事务外,R2)→
|
||
单原子事务(晋升 + digest + 裁决留痕)→ 一次 commit → 登记**验收后伏笔到期扫描**
|
||
BackgroundTask(M3-b/d,§6.2)。
|
||
|
||
到期扫描接 §5.5 步骤 3 的 `TODO(M3)` 占位:验收提交成功**后**,把
|
||
`current_ch > expected_close_to AND status≠CLOSED` 的伏笔确定性置 OVERDUE。
|
||
伏笔登记/回收(PARTIAL/CLOSED)的「采纳建议」经作者经 foreshadow 端点显式确认,
|
||
不在验收事务里 AI 静默改库(不变量 #3)。**§7.4 持久性局限**:BackgroundTask 进程内
|
||
跑、重启丢失——原型可接受(见 `services/foreshadow_scan.py`,不引入 jobs 表)。
|
||
"""
|
||
request_id = getattr(request.state, "request_id", None)
|
||
# R3:审稿真相从领域表重读(最近一条 chapter_reviews),不依赖 checkpoint。
|
||
history = await review_repo.list_for_chapter(project_id, chapter_no)
|
||
latest_review = history[0] if history else None
|
||
|
||
# 冲突 gate(R5,事务前拦截,不写库)。
|
||
assert_conflicts_resolved(latest_review, body.decisions)
|
||
|
||
log.info(
|
||
"accept_start",
|
||
project_id=str(project_id),
|
||
chapter_no=chapter_no,
|
||
request_id=request_id,
|
||
final_text_len=len(body.final_text),
|
||
decision_count=len(body.decisions),
|
||
has_review=latest_review is not None,
|
||
)
|
||
|
||
# R2:终稿 digest 提炼在**事务外**做(别在持开事务里跨网络调 LLM)。
|
||
digest_facts = await extract_digest_facts(
|
||
gateway,
|
||
final_text=body.final_text,
|
||
user_id=STUB_OWNER_ID,
|
||
project_id=project_id,
|
||
chapter_no=chapter_no,
|
||
)
|
||
|
||
outcome: AcceptOutcome = await run_accept_transaction(
|
||
session=session,
|
||
chapter_repo=chapter_repo,
|
||
digest_repo=digest_repo,
|
||
review_repo=review_repo,
|
||
project_id=project_id,
|
||
chapter_no=chapter_no,
|
||
final_text=body.final_text,
|
||
digest_facts=digest_facts,
|
||
latest_review=latest_review,
|
||
decisions=body.decisions,
|
||
)
|
||
|
||
# 验收**提交成功后**登记伏笔到期扫描 BackgroundTask(M3-b/d)。current_chapter=刚验收章号;
|
||
# 任务在请求 session 关闭后跑,故 run_overdue_scan 自建新 session 并 commit(§7.4 局限可接受)。
|
||
background_tasks.add_task(
|
||
run_overdue_scan,
|
||
session_factory,
|
||
project_id=project_id,
|
||
chapter_no=chapter_no,
|
||
request_id=request_id,
|
||
)
|
||
|
||
return AcceptResponse(
|
||
project_id=project_id,
|
||
chapter_no=chapter_no,
|
||
accepted_version=outcome.accepted_version,
|
||
digest_added=outcome.digest_added,
|
||
decisions_recorded=outcome.decisions_recorded,
|
||
review_id=outcome.review_id,
|
||
)
|