274 lines
11 KiB
Python
274 lines
11 KiB
Python
"""CRUD API for roundtable-planning drafts + recommendation history.
|
||
|
||
A draft is one full roundtable session (Step 1 intent + Step 2 multi-agent
|
||
roundtable). Bound to the requesting user (not the browser), so it follows the
|
||
account across devices/sessions — replacing the old ``localStorage`` storage.
|
||
|
||
This store is **strictly per-user**. taskId 深链聊天记录 (无界嵌入抽屉) live in a
|
||
separate, unpartitioned store — see ``/api/roundtable-task-drafts``
|
||
(``roundtable_task_drafts`` table) — so the two never mix.
|
||
|
||
Recommendation history is bound to a draft: the recommend dialog shows "last
|
||
recommendation for this session" on open, while still letting the user
|
||
re-analyze (which appends a new history row).
|
||
|
||
Routes (prefix ``/api/roundtable-drafts``):
|
||
GET "" list current user's drafts (lightweight)
|
||
POST "" create a draft
|
||
GET "/{draft_id}" full draft (incl. step1 / step2 blobs)
|
||
PUT "/{draft_id}" partial update (title / step / snapshots)
|
||
DELETE "/{draft_id}" delete
|
||
GET "/{draft_id}/recommendations" list recommendation history
|
||
POST "/{draft_id}/recommendations" append one recommendation run
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
from datetime import datetime
|
||
from typing import Any
|
||
from uuid import uuid4
|
||
|
||
from fastapi import APIRouter, HTTPException, Request
|
||
from pydantic import BaseModel, Field
|
||
|
||
from deerflow.persistence.roundtable_drafts.sql import DraftConcurrentWriteError
|
||
from deerflow.runtime.user_context import get_effective_user_id
|
||
|
||
logger = logging.getLogger(__name__)
|
||
router = APIRouter(prefix="/api/roundtable-drafts", tags=["roundtable-drafts"])
|
||
|
||
|
||
# ── schemas ────────────────────────────────────────────────────────────────
|
||
|
||
|
||
class DraftMetaResponse(BaseModel):
|
||
id: str
|
||
# 外部任务 id(无界嵌入抽屉传入);非空时草稿按 task_id 共享、不按 user 分权。
|
||
task_id: str | None = None
|
||
title: str = ""
|
||
furthest_step: int = 1
|
||
# 乐观锁版本号:PUT 时经 expected_version 回传,冲突 → 409 + current_version。
|
||
version: int = 0
|
||
# 本次会商所用业务链名(历史卡片展示);无链条(自由/recommend 模式)时为 None。
|
||
chain_title: str | None = None
|
||
created_at: datetime | str | None = None
|
||
updated_at: datetime | str | None = None
|
||
# 最近一次后台作业的状态(Phase 6 历史下拉「进行中」loading 用);无作业时为 None。
|
||
status: str | None = None
|
||
job_id: str | None = None
|
||
|
||
|
||
class DraftListResponse(BaseModel):
|
||
drafts: list[DraftMetaResponse]
|
||
|
||
|
||
class DraftResponse(DraftMetaResponse):
|
||
step1: dict[str, Any] | None = None
|
||
step2: dict[str, Any] | None = None
|
||
step3: dict[str, Any] | None = None
|
||
|
||
|
||
class DraftCreateRequest(BaseModel):
|
||
# Optional client-provided id so the browser can keep one stable id across
|
||
# the create + subsequent updates without an extra round-trip.
|
||
id: str | None = Field(default=None, max_length=64)
|
||
# 外部任务 id(无界嵌入抽屉);设置后草稿按 task_id 共享、不按 user 分权。
|
||
task_id: str | None = Field(default=None, max_length=128)
|
||
title: str = Field(default="", max_length=512)
|
||
furthest_step: int = Field(default=1, ge=1, le=3)
|
||
step1: dict[str, Any] | None = None
|
||
step2: dict[str, Any] | None = None
|
||
step3: dict[str, Any] | None = None
|
||
|
||
|
||
class DraftUpdateRequest(BaseModel):
|
||
# All optional → callers send only what changed. Auto-save omits ``title``
|
||
# so it never clobbers a manual rename.
|
||
title: str | None = Field(default=None, max_length=512)
|
||
furthest_step: int | None = Field(default=None, ge=1, le=3)
|
||
task_id: str | None = Field(default=None, max_length=128)
|
||
step1: dict[str, Any] | None = None
|
||
step2: dict[str, Any] | None = None
|
||
step3: dict[str, Any] | None = None
|
||
# 乐观锁:客户端基于的草稿版本号。缺省(旧客户端)= 盲写兼容;
|
||
# 提供时版本不匹配 → 409,body 带 current_version 供重拉合并。
|
||
expected_version: int | None = Field(default=None, ge=0)
|
||
|
||
|
||
class RecommendPickSchema(BaseModel):
|
||
agent_id: str
|
||
reason: str = ""
|
||
|
||
|
||
class RecommendCandidateSchema(BaseModel):
|
||
agent_id: str
|
||
name: str = ""
|
||
description: str = ""
|
||
|
||
|
||
class RecommendHistoryResponse(BaseModel):
|
||
id: str
|
||
draft_id: str
|
||
objective: str = ""
|
||
status: str = "done"
|
||
model: str | None = None
|
||
rationale: str | None = None
|
||
picks: list[RecommendPickSchema] = []
|
||
candidates: list[RecommendCandidateSchema] = []
|
||
created_at: datetime | str | None = None
|
||
|
||
|
||
class RecommendHistoryListResponse(BaseModel):
|
||
history: list[RecommendHistoryResponse]
|
||
|
||
|
||
class RecommendHistoryCreateRequest(BaseModel):
|
||
objective: str = Field(default="", max_length=512)
|
||
status: str = Field(default="done", max_length=32)
|
||
model: str | None = None
|
||
rationale: str | None = None
|
||
picks: list[RecommendPickSchema] = []
|
||
candidates: list[RecommendCandidateSchema] = []
|
||
|
||
|
||
# ── helpers ──────────────────────────────────────────────────────────────
|
||
|
||
|
||
def _current_user_id(request: Request) -> str:
|
||
user = getattr(request.state, "user", None)
|
||
if user is not None:
|
||
return str(user.id)
|
||
return get_effective_user_id()
|
||
|
||
|
||
def _get_store(request: Request):
|
||
store = getattr(request.app.state, "roundtable_draft_store", None)
|
||
if store is None:
|
||
raise HTTPException(status_code=503, detail="Roundtable draft store not available")
|
||
return store
|
||
|
||
|
||
# ── draft routes ───────────────────────────────────────────────────────────
|
||
|
||
|
||
@router.get("", response_model=DraftListResponse)
|
||
async def list_drafts(request: Request) -> DraftListResponse:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
rows = await store.list_drafts(user_id)
|
||
# 富化:给每个草稿带上「最近一次后台作业」的状态 + jobId(Phase 6 历史下拉 loading 用)。
|
||
# job store 缺失(DB 未配置)或查询失败都不阻断列表,静默退化为无状态。
|
||
job_store = getattr(request.app.state, "roundtable_job_store", None)
|
||
if job_store is not None:
|
||
try:
|
||
status_map = await job_store.status_map_by_user(user_id)
|
||
except Exception: # noqa: BLE001 — 富化失败不能拖垮列表
|
||
logger.warning("roundtable draft list: job status enrichment failed", exc_info=True)
|
||
status_map = {}
|
||
for r in rows:
|
||
info = status_map.get(r.get("id"))
|
||
if info:
|
||
r["status"] = info.get("status")
|
||
r["job_id"] = info.get("jobId")
|
||
return DraftListResponse(drafts=[DraftMetaResponse(**r) for r in rows])
|
||
|
||
|
||
@router.post("", response_model=DraftResponse, status_code=201)
|
||
async def create_draft(request: Request, body: DraftCreateRequest) -> DraftResponse:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
row = await store.create_draft(
|
||
user_id,
|
||
{
|
||
"id": (body.id or uuid4().hex),
|
||
"task_id": (body.task_id or None),
|
||
"title": body.title.strip(),
|
||
"furthest_step": body.furthest_step,
|
||
"step1": body.step1,
|
||
"step2": body.step2,
|
||
"step3": body.step3,
|
||
},
|
||
)
|
||
return DraftResponse(**row)
|
||
|
||
|
||
@router.get("/{draft_id}", response_model=DraftResponse)
|
||
async def get_draft(request: Request, draft_id: str) -> DraftResponse:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
row = await store.get_draft(draft_id, user_id)
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="Draft not found")
|
||
return DraftResponse(**row)
|
||
|
||
|
||
@router.put("/{draft_id}", response_model=DraftResponse)
|
||
async def update_draft(request: Request, draft_id: str, body: DraftUpdateRequest) -> DraftResponse:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
# Build kwargs that distinguish "omitted" from "explicitly None". The store
|
||
# uses a sentinel; here we pass only the fields the client actually sent.
|
||
fields = body.model_dump(exclude_unset=True)
|
||
try:
|
||
row = await store.update_draft(draft_id, user_id, **fields)
|
||
except DraftConcurrentWriteError as err:
|
||
# 版本冲突:detail 里带 current_version,前端据此重拉最新草稿、本地合并后重试。
|
||
raise HTTPException(
|
||
status_code=409,
|
||
detail={
|
||
"error": "draft_version_conflict",
|
||
"message": f"草稿已被并发修改(expected={err.expected}, actual={err.actual})",
|
||
"current_version": err.actual,
|
||
},
|
||
) from err
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="Draft not found")
|
||
return DraftResponse(**row)
|
||
|
||
|
||
@router.delete("/{draft_id}", status_code=204)
|
||
async def delete_draft(request: Request, draft_id: str) -> None:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
deleted = await store.delete_draft(draft_id, user_id)
|
||
if not deleted:
|
||
raise HTTPException(status_code=404, detail="Draft not found")
|
||
|
||
|
||
# ── recommendation history routes ────────────────────────────────────────
|
||
|
||
|
||
@router.get("/{draft_id}/recommendations", response_model=RecommendHistoryListResponse)
|
||
async def list_recommendations(request: Request, draft_id: str) -> RecommendHistoryListResponse:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
rows = await store.list_recommendations(draft_id, user_id)
|
||
return RecommendHistoryListResponse(history=[RecommendHistoryResponse(**r) for r in rows])
|
||
|
||
|
||
@router.post("/{draft_id}/recommendations", response_model=RecommendHistoryResponse, status_code=201)
|
||
async def add_recommendation(
|
||
request: Request, draft_id: str, body: RecommendHistoryCreateRequest
|
||
) -> RecommendHistoryResponse:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
# Ownership: the draft must exist and belong to the caller.
|
||
draft = await store.get_draft(draft_id, user_id)
|
||
if draft is None:
|
||
raise HTTPException(status_code=404, detail="Draft not found")
|
||
row = await store.add_recommendation(
|
||
user_id,
|
||
draft_id,
|
||
{
|
||
"id": uuid4().hex,
|
||
"objective": body.objective,
|
||
"status": body.status,
|
||
"model": body.model,
|
||
"rationale": body.rationale,
|
||
"picks": [p.model_dump() for p in body.picks],
|
||
"candidates": [c.model_dump() for c in body.candidates],
|
||
},
|
||
)
|
||
return RecommendHistoryResponse(**row)
|