draft/rewrite/review 三个流式端点返回 StreamingResponse,OpenAPI 原按默认 application/json 标注 200,schema 不诚实。新增共享 _SSE_RESPONSE 声明 200 的 content 为 text/event-stream 并挂到三个装饰器。export_openapi 通过,schema 现 出现 3 处 text/event-stream。未跑前端 gen:api(按要求)。
730 lines
29 KiB
Python
730 lines
29 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, Any
|
||
|
||
from fastapi import APIRouter, BackgroundTasks, Depends, Request
|
||
from fastapi.responses import StreamingResponse
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
from ww_agents import ClarifyDecision, clarify_rewrite_spec
|
||
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,
|
||
error_event,
|
||
normalize_deltas,
|
||
normalize_review,
|
||
run_clarify,
|
||
stream_chapter_draft,
|
||
stream_chapter_rewrite,
|
||
)
|
||
from ww_db import get_session
|
||
from ww_llm_gateway import Gateway
|
||
from ww_shared import AppError, ErrorCode, ErrorEnvelope
|
||
|
||
from ww_api.logging_config import get_logger
|
||
from ww_api.pagination import PageDep
|
||
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,
|
||
ReviewConflictView,
|
||
ReviewHistoryItem,
|
||
ReviewHistoryResponse,
|
||
ReviewRequest,
|
||
RewriteClarifyRequest,
|
||
RewriteStreamRequest,
|
||
)
|
||
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_clarify_gateway,
|
||
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"])
|
||
|
||
# OpenAPI 错误响应声明(让 TS 客户端拿到类型化错误形,§7.1)。
|
||
_NOT_FOUND: dict[int | str, dict[str, Any]] = {
|
||
404: {"model": ErrorEnvelope, "description": "资源不存在"}
|
||
}
|
||
_ACCEPT_ERRORS: dict[int | str, dict[str, Any]] = {
|
||
404: {"model": ErrorEnvelope, "description": "资源不存在"},
|
||
409: {"model": ErrorEnvelope, "description": "存在未裁决冲突"},
|
||
503: {"model": ErrorEnvelope, "description": "LLM 不可用"},
|
||
}
|
||
# SSE 端点:如实声明 200 的 media type 为 text/event-stream(否则 OpenAPI 误标 json)。
|
||
_SSE_RESPONSE: dict[int | str, dict[str, Any]] = {
|
||
200: {
|
||
"description": "SSE 事件流(`event: <name>\\ndata: <json>`,见 memory/contracts.md)",
|
||
"content": {"text/event-stream": {}},
|
||
}
|
||
}
|
||
|
||
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)]
|
||
ClarifyGatewayDep = Annotated[Gateway, Depends(get_clarify_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,
|
||
tone=body.tone,
|
||
ending_type=body.ending_type,
|
||
narrative_pov=body.narrative_pov,
|
||
),
|
||
)
|
||
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, page: PageDep) -> ProjectListResponse:
|
||
"""作品列表(按 created_at 升序,分页)。`?limit=&offset=`,默认第一页。"""
|
||
views = await repo.list_for_owner(STUB_OWNER_ID, limit=page.limit, offset=page.offset)
|
||
return ProjectListResponse(projects=[_to_response(v) for v in views])
|
||
|
||
|
||
@router.get("/{project_id}", responses=_NOT_FOUND)
|
||
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", responses=_NOT_FOUND)
|
||
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", responses=_NOT_FOUND)
|
||
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", responses=_SSE_RESPONSE)
|
||
async def stream_draft(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
request: Request,
|
||
repos: MemoryReposDep,
|
||
gateway: GatewayDep,
|
||
injection_repo: InjectionRepoDep,
|
||
project_repo: ProjectRepoDep,
|
||
session: Annotated[AsyncSession, Depends(get_session)],
|
||
body: DraftStreamRequest | None = None,
|
||
) -> StreamingResponse:
|
||
"""流式写章草稿:组装记忆(含作者注入覆盖 + 本章指令)→ 网关流 → 归一 SSE → event-stream。
|
||
|
||
`body.directive`(可选,T4-b)是临时本章指令,直通 assemble→volatile(不持久化);
|
||
无 body 的旧调用方仍可用(向后兼容)。
|
||
|
||
项目不存在 → 404(在触网关前 fail-fast):否则非法 project_id 会静默烧一次付费/限流的
|
||
LLM 调用并返回 200(QA C1)。
|
||
"""
|
||
if await project_repo.get(STUB_OWNER_ID, project_id) is None:
|
||
raise AppError(ErrorCode.NOT_FOUND, f"project not found: {project_id}")
|
||
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 暴露)。
|
||
# 尾部 commit 包 try/except:失败则记 sse_commit_failed(账本静默丢失须可查,P0-4)。
|
||
try:
|
||
await session.commit()
|
||
except Exception: # noqa: BLE001 — 流已发完,commit 失败不能再改响应;至少记错误。
|
||
log.error("sse_commit_failed", request_id=request_id, endpoint="draft")
|
||
|
||
return StreamingResponse(
|
||
_frames(),
|
||
media_type="text/event-stream",
|
||
headers={"cache-control": "no-cache", "x-accel-buffering": "no"},
|
||
)
|
||
|
||
|
||
@router.post("/{project_id}/chapters/{chapter_no}/rewrite", responses=_SSE_RESPONSE)
|
||
async def stream_rewrite(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
request: Request,
|
||
repos: MemoryReposDep,
|
||
gateway: GatewayDep,
|
||
injection_repo: InjectionRepoDep,
|
||
project_repo: ProjectRepoDep,
|
||
session: Annotated[AsyncSession, Depends(get_session)],
|
||
body: RewriteStreamRequest,
|
||
) -> StreamingResponse:
|
||
"""整章再沟通/重写:组装记忆 + 当前草稿 + 作者意见 → 网关流 → 归一 SSE → event-stream。
|
||
|
||
HITL:新版停前端/草稿,接受才落库(不变量 #3);rewrite 是工具非写节点,不破坏章
|
||
纯函数性(不变量 #7)。项目不存在 → 404(触网关前 fail-fast,同 draft,QA C1)。
|
||
"""
|
||
if await project_repo.get(STUB_OWNER_ID, project_id) is None:
|
||
raise AppError(ErrorCode.NOT_FOUND, f"project not found: {project_id}")
|
||
request_id = getattr(request.state, "request_id", None)
|
||
override = await injection_repo.get(project_id, chapter_no)
|
||
context = await assemble(repos, project_id, chapter_no, override=override)
|
||
log.info(
|
||
"rewrite_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),
|
||
prior_draft_len=len(body.prior_draft),
|
||
feedback_len=len(body.feedback),
|
||
)
|
||
|
||
deltas = stream_chapter_rewrite(
|
||
gateway,
|
||
stable_core=context.stable_core,
|
||
volatile=context.volatile,
|
||
prior_draft=body.prior_draft,
|
||
feedback=body.feedback,
|
||
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)
|
||
try:
|
||
await session.commit()
|
||
except Exception: # noqa: BLE001 — 流已发完,commit 失败不能再改响应;至少记错误。
|
||
log.error("sse_commit_failed", request_id=request_id, endpoint="rewrite")
|
||
|
||
return StreamingResponse(
|
||
_frames(),
|
||
media_type="text/event-stream",
|
||
headers={"cache-control": "no-cache", "x-accel-buffering": "no"},
|
||
)
|
||
|
||
|
||
# 整章重写预检的草稿摘录上界:只喂 analyst 预检**开头一小段**供锚定,绝不整章入 LLM
|
||
# (守成本/上下文,且预检只需判意见含糊度、不需读全章)。
|
||
_REWRITE_CLARIFY_EXCERPT_MAX = 3000
|
||
|
||
|
||
def _build_rewrite_clarify_input(feedback: str, prior_draft: str | None) -> str:
|
||
"""组整章重写预检输入:作者章级意见 + 当前章草稿的有界摘录(分块标注,供教条区分)。
|
||
|
||
只取 `prior_draft` 开头 `_REWRITE_CLARIFY_EXCERPT_MAX` 字符作锚定摘录——预检不需读全章,
|
||
避免把 200k 整章喂给 analyst 预检(成本/上下文/脱敏)。无草稿时只喂意见(退化,仍可判含糊)。
|
||
"""
|
||
parts = [f"【作者的整章修改意见】\n{feedback}"]
|
||
if prior_draft:
|
||
excerpt = prior_draft[:_REWRITE_CLARIFY_EXCERPT_MAX]
|
||
parts.append(f"【当前章节草稿(摘录,仅供锚定)】\n{excerpt}")
|
||
return "\n\n".join(parts)
|
||
|
||
|
||
@router.post(
|
||
"/{project_id}/chapters/{chapter_no}/rewrite/clarify",
|
||
responses={
|
||
404: {"model": ErrorEnvelope, "description": "项目不存在"},
|
||
503: {"model": ErrorEnvelope, "description": "LLM 不可用"},
|
||
},
|
||
)
|
||
async def clarify_rewrite(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
body: RewriteClarifyRequest,
|
||
request: Request,
|
||
project_repo: ProjectRepoDep,
|
||
gateway: ClarifyGatewayDep,
|
||
session: Annotated[AsyncSession, Depends(get_session)],
|
||
) -> ClarifyDecision:
|
||
"""整章重写预检澄清(WFW-9 M2 路线A 两阶段之「问题」阶段):非流式 JSON,只判不改。
|
||
|
||
镜像 refine 侧 clarify——只判「作者这条整章重写意见清不清楚」,含糊则反问 1 问 + 给选项,
|
||
清楚则给确认语放行;**绝不改正文**(整章重写仍走既有 rewrite SSE 端点)。analyst 档
|
||
(不变量 #2)。只取 `prior_draft` 有界摘录喂预检(不整章入 LLM)。
|
||
|
||
项目不存在 → 404(触网关前 fail-fast,同 rewrite/draft);无凭据 → 503(dep 解析阶段拦下)。
|
||
**只读不写业务表**(不变量 #3);末尾 `commit()` 让网关 usage_ledger 落库(add-only,否则
|
||
记账静默丢失,同 rewrite 纪律)。判别/校验失败在 `run_clarify` 内确定性回退
|
||
`need_clarification=false`(不阻塞整章重写主链路)。
|
||
"""
|
||
request_id = getattr(request.state, "request_id", None)
|
||
|
||
if await project_repo.get(STUB_OWNER_ID, project_id) is None:
|
||
raise AppError(ErrorCode.NOT_FOUND, f"project not found: {project_id}")
|
||
|
||
context = _build_rewrite_clarify_input(body.feedback, body.prior_draft)
|
||
decision = await run_clarify(
|
||
clarify_rewrite_spec, # 不变量 #2:analyst 档,不传 model
|
||
context=context,
|
||
gateway=gateway,
|
||
user_id=STUB_OWNER_ID,
|
||
project_id=project_id,
|
||
)
|
||
|
||
# 提交边界:只读不写业务表,但网关 ledger add-only → 端点末尾 commit(否则记账丢失)。
|
||
await session.commit()
|
||
|
||
log.info(
|
||
"rewrite_clarify_done",
|
||
project_id=str(project_id),
|
||
chapter_no=chapter_no,
|
||
request_id=request_id,
|
||
feedback_len=len(body.feedback), # 脱敏:只记长度,不记正文(不变量/日志纪律)
|
||
need_clarification=decision.need_clarification,
|
||
question_count=len(decision.questions),
|
||
)
|
||
return decision
|
||
|
||
|
||
@router.put("/{project_id}/chapters/{chapter_no}/draft")
|
||
async def save_draft(
|
||
project_id: uuid.UUID,
|
||
chapter_no: int,
|
||
body: DraftSaveRequest,
|
||
repo: ChapterRepoDep,
|
||
session: Annotated[AsyncSession, Depends(get_session)],
|
||
) -> DraftResponse:
|
||
"""自动保存:幂等 upsert 草稿(同章节覆盖同一行,版次不爆炸)。"""
|
||
view = await repo.save_draft(project_id, chapter_no, text=body.text)
|
||
# repo.save_draft 只 flush,端点统一提交。
|
||
await session.commit()
|
||
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", responses=_NOT_FOUND)
|
||
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", responses=_SSE_RESPONSE)
|
||
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]:
|
||
# graph.ainvoke 包 try/except:抛错则先发 error 事件再 return(否则流被截断、
|
||
# 客户端收不到 error,P0-4)。AppError 用其 code/message;其余归一为 INTERNAL(不泄异常)。
|
||
try:
|
||
final = await graph.ainvoke(initial)
|
||
except AppError as exc:
|
||
log.warning("review_stream_error", code=str(exc.code), request_id=request_id)
|
||
yield _encode_sse(
|
||
error_event(code=str(exc.code), message=exc.message, request_id=request_id)
|
||
)
|
||
return
|
||
except Exception as exc: # noqa: BLE001 — 边界兜底:任何意外归一为 error 事件,不泄异常
|
||
log.error("review_stream_unexpected_error", error=str(exc), request_id=request_id)
|
||
yield _encode_sse(
|
||
error_event(
|
||
code=str(ErrorCode.INTERNAL),
|
||
message="internal error during review",
|
||
request_id=request_id,
|
||
)
|
||
)
|
||
return
|
||
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 端点)。
|
||
# 尾部 commit 包 try/except:失败则记 sse_commit_failed(P0-4)。
|
||
try:
|
||
await session.commit()
|
||
except Exception: # noqa: BLE001 — 流已发完,commit 失败不能再改响应;至少记错误。
|
||
log.error("sse_commit_failed", request_id=request_id, endpoint="review")
|
||
|
||
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=[ReviewConflictView.model_validate(c) for c in v.conflicts],
|
||
foreshadow_sug=v.foreshadow_sug,
|
||
style=v.style,
|
||
pace=v.pace,
|
||
characterization=v.characterization,
|
||
health_score=v.health_score,
|
||
decisions=v.decisions,
|
||
)
|
||
for v in views
|
||
]
|
||
return ReviewHistoryResponse(reviews=items)
|
||
|
||
|
||
@router.post("/{project_id}/chapters/{chapter_no}/accept", responses=_ACCEPT_ERRORS)
|
||
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,
|
||
)
|