Files
writer-work-flow/apps/api/ww_api/routers/projects.py
Yaojia Wang b523b4fd21 feat: M1 — 立项→写章草稿(SSE)→自动保存;连一家 provider
- 薄自建 LLM 网关:OpenAI 兼容适配器(DeepSeek) + instructor 结构化输出 + usage_ledger 记账 + 档位路由
- 记忆服务 assemble:确定性选择(显式+主角+近况) + 渲染卡 + 缓存断点(中性文本)
- LangGraph 写章节点 + Postgres checkpointer + SSE 归一(token/done/error)
- API:立项 + 写章 draft(SSE) + PUT 自动保存 + 提供商凭据(Fernet 加密/测试连接)
- 前端:AppShell + 作品库 + 5 步立项向导 + 写作工作台(流式打字机+自动保存) + 设置页
- M1 E2E:真实 DB + mock 网关零 token 走通闭环
2026-06-18 11:38:28 +02:00

361 lines
13 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.

"""项目(立项)+ 章节草稿端点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 流式写章草稿SSEtext/event-stream
- PUT /projects/:id/chapters/:no/draft 自动保存草稿(幂等 upsert
不变量writer 只产草稿(不提升 accepted 版本/不抽 digest/不写状态——属 M2 验收);
agent 只传 tierDB 按 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, 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.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.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.projects import (
AcceptRequest,
AcceptResponse,
DraftResponse,
DraftSaveRequest,
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.project_deps import (
get_chapter_repo,
get_digest_append_repo,
get_digest_gateway,
get_memory_repos,
get_project_repo,
get_review_gateway,
get_review_repo,
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)]
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)
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,
session: Annotated[AsyncSession, Depends(get_session)],
) -> StreamingResponse:
"""流式写章草稿:组装记忆 → 网关流 → 归一为 SSE 事件 → text/event-stream。"""
request_id = getattr(request.state, "request_id", None)
context = await assemble(repos, project_id, chapter_no)
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),
)
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),
)
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,
chapter_repo: ChapterRepoDep,
digest_repo: DigestRepoDep,
review_repo: ReviewRepoDep,
gateway: DigestGatewayDep,
session: Annotated[AsyncSession, Depends(get_session)],
) -> AcceptResponse:
"""验收事务 + 冲突 gate§5.5gate事务前→ 终稿提炼 digest事务外R2
单原子事务(晋升 + digest + 裁决留痕)→ 一次 commit。
"""
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
# 冲突 gateR5事务前拦截不写库
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,
)
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,
)