"""Gateway API for the AgentScope report-collaboration workbench. Plan/run writes enqueue durable jobs and return immediately. AgentScope execution is not started from this router. ``GET /runs/{id}/stream`` is a read-only durable SSE tail: disconnect does not cancel the run. """ from __future__ import annotations from typing import Any from fastapi import APIRouter, Header, HTTPException, Query, Request from fastapi.responses import JSONResponse, Response, StreamingResponse from pydantic import BaseModel, Field from app.gateway.deps import get_current_user from app.report_collaboration.agentscope_runtime.model_resolver import freeze_plan_models from app.report_collaboration.contracts.snapshots import SessionSnapshot from app.report_collaboration.execution.live_hub import iter_run_sse, replay_cursor from app.report_collaboration.interventions.intent_router import build_runtime_intent_router from app.report_collaboration.interventions.service import RuntimeCommandService from app.report_collaboration.planning import PlanProposalService from app.report_collaboration.reporting import ConservativeRewriteGenerator, ModelRewriteGenerator, ReportRewriteService, freeze_report_template from app.report_collaboration.reporting.templates import requirement_snapshot_from_record from app.report_collaboration.requirement import RequirementResolver from app.report_collaboration.security import freeze_run_budget from deerflow.config.app_config import get_app_config from deerflow.persistence.report_collaboration import ( ReportCollaborationConflictError, ReportCollaborationLimitError, ReportCollaborationNotFoundError, ReportCollaborationStore, ReportCollaborationValidationError, ) router = APIRouter(prefix="/api/report-collaboration", tags=["report-collaboration"]) class CreateSessionBody(BaseModel): title: str | None = Field(default=None, max_length=255) class UpdateSessionBody(BaseModel): title: str | None = Field(default=None, max_length=255) class SendMessageBody(BaseModel): text: str = Field(..., min_length=1, max_length=20_000) target_node_id: str | None = None agent_run_id: str | None = None class StartRunBody(BaseModel): plan_id: str = Field(..., min_length=1, max_length=64) class CreateCommandBody(BaseModel): text: str = Field(..., min_length=1, max_length=20_000) target_node_id: str | None = None agent_run_id: str | None = None class CreateRewriteBody(BaseModel): instruction: str = Field(..., min_length=1, max_length=20_000) section: str | None = None def _store(request: Request) -> ReportCollaborationStore: store = getattr(request.app.state, "report_collaboration_store", None) if store is None: raise HTTPException(status_code=503, detail="Report collaboration store is unavailable") return store def _runtime_commands(request: Request) -> RuntimeCommandService: service = getattr(request.app.state, "report_collaboration_command_service", None) if service is None: app_config = get_app_config() service = RuntimeCommandService( _store(request), router=build_runtime_intent_router(app_config), config=app_config.report_collaboration, ) request.app.state.report_collaboration_command_service = service return service def _require_enabled() -> JSONResponse | None: if not get_app_config().report_collaboration.enabled: return _error(503, "REPORT_COLLABORATION_DISABLED", "报告协作工作台未启用") return None async def _require_user(request: Request) -> str: user_id = await get_current_user(request) if not user_id: raise HTTPException(status_code=401, detail="Authentication required") return str(user_id) def _error(status: int, code: str, message: str, detail: Any = None) -> JSONResponse: body: dict[str, Any] = {"code": code, "message": message} if detail is not None: body["detail"] = detail return JSONResponse(status_code=status, content=body) def _require_idempotency(idempotency_key: str | None) -> str | JSONResponse: key = (idempotency_key or "").strip() if not key: return _error(422, "IDEMPOTENCY_KEY_REQUIRED", "写操作必须携带 X-Idempotency-Key") if len(key) > 128: return _error(422, "IDEMPOTENCY_KEY_INVALID", "X-Idempotency-Key 过长") return key def _map_store_error(exc: Exception) -> JSONResponse: if isinstance(exc, ReportCollaborationNotFoundError): return _error(404, "NOT_FOUND", str(exc)) if isinstance(exc, ReportCollaborationConflictError): return _error(409, exc.code, str(exc), detail={"current_revision": exc.current_revision}) if isinstance(exc, ReportCollaborationValidationError): return _error(422, exc.code, str(exc)) if isinstance(exc, ReportCollaborationLimitError): return _error(429, exc.code, str(exc), detail={"usage": exc.usage, "limits": exc.limits}) raise exc def _rewrite_service(store: ReportCollaborationStore) -> ReportRewriteService: cfg = get_app_config().report_collaboration generator = ModelRewriteGenerator(cfg.writer_model) if cfg.writer_model else ConservativeRewriteGenerator() return ReportRewriteService(store, generator=generator) async def _frozen_run_template(request: Request, session: dict[str, Any], plan: dict[str, Any]) -> dict[str, Any]: requirement = requirement_snapshot_from_record(session) structure_id = plan.get("report_structure_id") or (requirement.report_structure_id if requirement else None) structure = None structure_store = getattr(request.app.state, "report_structure_store", None) if structure_id and structure_store is not None: structure = await structure_store.get_structure(str(structure_id)) return freeze_report_template(structure, requirement=requirement).model_dump(mode="json") async def _owned_session(store: ReportCollaborationStore, session_id: str, user_id: str) -> dict[str, Any] | JSONResponse: record = await store.get_session_record(session_id) if record is None or record.get("owner_id") != user_id: return _error(404, "NOT_FOUND", f"report collaboration session not found: {session_id}") return record async def _owned_run(store: ReportCollaborationStore, run_id: str, user_id: str) -> dict[str, Any] | JSONResponse: run = await store.get_run(run_id) if run is None: return _error(404, "NOT_FOUND", f"report collaboration run not found: {run_id}") owned = await _owned_session(store, run["session_id"], user_id) if isinstance(owned, JSONResponse): return owned return run @router.post("/sessions") async def create_session( request: Request, body: CreateSessionBody | None = None, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key payload = body or CreateSessionBody() try: return await _store(request).create_session(owner_id=user_id, title=payload.title or "", idempotency_key=key) except Exception as exc: return _map_store_error(exc) @router.get("/sessions") async def list_sessions(request: Request): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) sessions = await _store(request).list_sessions(user_id) return {"sessions": sessions} @router.get("/sessions/{session_id}") async def get_session_snapshot(request: Request, session_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned snapshot = await store.get_snapshot(session_id) return SessionSnapshot.model_validate(snapshot).model_dump() @router.patch("/sessions/{session_id}") async def update_session( request: Request, session_id: str, body: UpdateSessionBody, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), x_expected_revision: int | None = Header(default=None, alias="X-Expected-Revision"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned try: updated = await store.update_session(session_id, title=body.title, expected_revision=x_expected_revision) await store.create_command(session_id, idempotency_key=key, text=body.title or "", operation="update_session", status="completed") return updated except Exception as exc: return _map_store_error(exc) @router.delete("/sessions/{session_id}") async def delete_session(request: Request, session_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned await store.delete_session(session_id) return Response(status_code=204) @router.get("/sessions/{session_id}/messages") async def list_messages(request: Request, session_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned return {"messages": await store.list_messages(session_id)} @router.post("/sessions/{session_id}/messages") async def send_message( request: Request, session_id: str, body: SendMessageBody, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned metadata: dict[str, Any] = {"idempotency_key": key} if body.target_node_id: metadata["target_node_id"] = body.target_node_id if body.agent_run_id: metadata["agent_run_id"] = body.agent_run_id try: existing = await store.get_message_by_idempotency(key) if existing is not None: return existing snapshot = await store.get_snapshot(session_id) active = snapshot.get("run") or {} if str(active.get("status") or "") in { "queued", "planning", "running", "awaiting_input", "reviewing", "recovering", }: return _error( 409, "RUN_ACTIVE", "运行进行中请使用 /runs/{run_id}/commands 发送干预,不要走规划前消息接口", detail={"run_id": active.get("id"), "status": active.get("status")}, ) created = await store.create_message( session_id, role="human", content=body.text, idempotency_key=key, metadata=metadata, agent_run_id=body.agent_run_id, node_run_id=body.target_node_id, ) await RequirementResolver(store).handle_user_message(session_id, user_message=created, idempotency_key=key) return created except Exception as exc: return _map_store_error(exc) @router.post("/sessions/{session_id}/plan-requests") async def request_plans( request: Request, session_id: str, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), x_expected_revision: int | None = Header(default=None, alias="X-Expected-Revision"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned if owned.get("selected_plan_id"): return _error(409, "PLAN_ALREADY_SELECTED", "会话已选定方案,无法重新规划") try: command = await store.enqueue_plan_request(session_id, idempotency_key=key, expected_revision=x_expected_revision) if command.get("status") == "pending": await PlanProposalService( store, agent_store=getattr(request.app.state, "agent_store", None), ).fulfill(session_id, user_id=user_id, command_id=command["command_id"], idempotency_key=key) return Response(status_code=204) except Exception as exc: return _map_store_error(exc) @router.get("/sessions/{session_id}/plans") async def list_plans(request: Request, session_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned return {"plans": await store.list_plans(session_id)} @router.post("/sessions/{session_id}/plans/{plan_id}/select") async def select_plan( request: Request, session_id: str, plan_id: str, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), x_expected_revision: int | None = Header(default=None, alias="X-Expected-Revision"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned try: result = await store.select_plan(session_id, plan_id, idempotency_key=key, expected_revision=x_expected_revision) return {"selected_plan_id": result["selected_plan_id"]} except Exception as exc: return _map_store_error(exc) @router.post("/sessions/{session_id}/runs") async def start_run( request: Request, session_id: str, body: StartRunBody, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), x_expected_revision: int | None = Header(default=None, alias="X-Expected-Revision"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned try: plans = await store.list_plans(session_id) plan = next((item for item in plans if item.get("id") == body.plan_id), None) if plan is None: return _error(404, "NOT_FOUND", f"report collaboration plan not found: {body.plan_id}") bundle = freeze_plan_models(plan, get_app_config()) cfg = get_app_config().report_collaboration run = await store.create_run( session_id, plan_id=body.plan_id, idempotency_key=key, expected_revision=x_expected_revision, role_snapshots=bundle.dump_by_role(), template_snapshot=await _frozen_run_template(request, owned, plan), budget_snapshot=freeze_run_budget(cfg), max_active_runs_per_owner=cfg.max_concurrent_runs_per_user, ) existing_events = await store.list_events(run["id"], after_seq=0, limit=1) if not existing_events: await store.append_event( session_id=session_id, run_id=run["id"], event_type="heartbeat", data={ "audit": "model_snapshot", "roles": [item.model_dump() for item in bundle.by_role.values()], "fallbacks": [item.model_dump() for item in bundle.fallbacks()], }, ) return run except Exception as exc: return _map_store_error(exc) @router.get("/sessions/{session_id}/reports") async def list_reports(request: Request, session_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned return {"versions": await store.list_report_versions(session_id)} @router.post("/sessions/{session_id}/report-rewrites") async def create_rewrite( request: Request, session_id: str, body: CreateRewriteBody, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned try: return await _rewrite_service(store).create_candidate( session_id, idempotency_key=key, instruction=body.instruction, section=body.section, ) except Exception as exc: return _map_store_error(exc) @router.post("/sessions/{session_id}/report-rewrites/{rewrite_id}/apply") async def apply_rewrite( request: Request, session_id: str, rewrite_id: str, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned try: version, _created = await _rewrite_service(store).apply_candidate(session_id, rewrite_id, idempotency_key=key) return version except Exception as exc: return _map_store_error(exc) @router.post("/sessions/{session_id}/report-versions/{version_id}/restore") async def restore_report( request: Request, session_id: str, version_id: str, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) owned = await _owned_session(store, session_id, user_id) if isinstance(owned, JSONResponse): return owned try: version, _created = await _rewrite_service(store).restore(session_id, version_id, idempotency_key=key) return version except Exception as exc: return _map_store_error(exc) @router.get("/runs/{run_id}") async def get_run(request: Request, run_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) run = await _owned_run(_store(request), run_id, user_id) if isinstance(run, JSONResponse): return run return run @router.get("/runs/{run_id}/events") async def list_run_events( request: Request, run_id: str, after_seq: int = Query(default=0, ge=0), limit: int = Query(default=200, ge=1, le=1000), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) run = await _owned_run(store, run_id, user_id) if isinstance(run, JSONResponse): return run return {"events": await store.list_events(run_id, after_seq=after_seq, limit=limit)} @router.get("/runs/{run_id}/stream") async def stream_run( request: Request, run_id: str, after_seq: int = Query(default=0, ge=0), last_event_id: str | None = Header(default=None, alias="Last-Event-ID"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) run = await _owned_run(store, run_id, user_id) if isinstance(run, JSONResponse): return run hub = getattr(request.app.state, "report_collaboration_live_hub", None) cursor = replay_cursor(last_event_id or request.headers.get("last-event-id"), after_seq) settings = getattr(request.app.state, "report_collaboration_sse_settings", None) async def generator(): async def _disconnected() -> bool: return await request.is_disconnected() async for chunk in iter_run_sse( store=store, hub=hub, run_id=run_id, cursor=cursor, is_disconnected=_disconnected, settings=settings, ): yield chunk return StreamingResponse( generator(), media_type="text/event-stream", headers={ "Cache-Control": "no-cache, no-transform", "Connection": "keep-alive", "X-Accel-Buffering": "no", "X-Report-Collaboration-Run-Status": str(run.get("status") or ""), }, ) @router.post("/runs/{run_id}/cancel") async def cancel_run( request: Request, run_id: str, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) run = await _owned_run(store, run_id, user_id) if isinstance(run, JSONResponse): return run try: return await store.request_cancel_run(run_id, idempotency_key=key) except Exception as exc: return _map_store_error(exc) @router.post("/runs/{run_id}/commands") async def create_command( request: Request, run_id: str, body: CreateCommandBody, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) run = await _owned_run(store, run_id, user_id) if isinstance(run, JSONResponse): return run try: return await _runtime_commands(request).submit( run_id, text=body.text, idempotency_key=key, target_node_id=body.target_node_id, agent_run_id=body.agent_run_id, ) except Exception as exc: return _map_store_error(exc) @router.post("/runs/{run_id}/commands/{command_id}/confirm") async def confirm_command( request: Request, run_id: str, command_id: str, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) run = await _owned_run(store, run_id, user_id) if isinstance(run, JSONResponse): return run try: return await _runtime_commands(request).confirm(run_id, command_id, idempotency_key=key) except Exception as exc: return _map_store_error(exc) @router.post("/runs/{run_id}/commands/{command_id}/cancel") async def cancel_command( request: Request, run_id: str, command_id: str, x_idempotency_key: str | None = Header(default=None, alias="X-Idempotency-Key"), ): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) key = _require_idempotency(x_idempotency_key) if isinstance(key, JSONResponse): return key store = _store(request) run = await _owned_run(store, run_id, user_id) if isinstance(run, JSONResponse): return run try: return await _runtime_commands(request).cancel(run_id, command_id, idempotency_key=key) except Exception as exc: return _map_store_error(exc) @router.get("/runs/{run_id}/nodes/{node_id}/contract") async def get_node_contract(request: Request, run_id: str, node_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) run = await _owned_run(store, run_id, user_id) if isinstance(run, JSONResponse): return run try: return await store.get_node_contract(run_id, node_id) except Exception as exc: return _map_store_error(exc) @router.get("/runs/{run_id}/artifacts") async def list_artifacts(request: Request, run_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) store = _store(request) run = await _owned_run(store, run_id, user_id) if isinstance(run, JSONResponse): return run return {"artifacts": await store.list_artifacts(run_id)} @router.get("/runs/{run_id}/sources") async def list_sources(request: Request, run_id: str): if disabled := _require_enabled(): return disabled user_id = await _require_user(request) run = await _owned_run(_store(request), run_id, user_id) if isinstance(run, JSONResponse): return run return {"sources": []}