1849 lines
78 KiB
Python
1849 lines
78 KiB
Python
"""API for the standalone position-collaboration workspace.
|
||
|
||
Unlike the existing multi-agent roundtable, this feature has no coordinator and
|
||
no background roundtable job. The API only persists the task/intent/chain
|
||
snapshot and indexes one direct-agent node per business-chain seat. LangGraph
|
||
remains the live execution source, while SQL also stores recovery copies of
|
||
all conversations, drafts and delivered artifact bodies so checkpoint-volume
|
||
replacement cannot turn an existing collaboration history into a blank page.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
from datetime import UTC, datetime
|
||
from pathlib import Path
|
||
from typing import Any
|
||
from uuid import uuid4
|
||
|
||
from fastapi import APIRouter, HTTPException, Query, Request
|
||
from pydantic import BaseModel, Field, field_validator
|
||
|
||
from app.gateway.path_utils import aresolve_thread_virtual_path
|
||
from deerflow.persistence.position_roundtable.delivery import (
|
||
artifact_signature as _artifact_signature,
|
||
)
|
||
from deerflow.persistence.position_roundtable.delivery import (
|
||
is_intent_position_node as _is_intent_position_by_id,
|
||
)
|
||
from deerflow.persistence.position_roundtable.delivery import (
|
||
is_json_artifact as _is_json_artifact,
|
||
)
|
||
from deerflow.persistence.position_roundtable.delivery import (
|
||
is_markdown_artifact as _is_markdown_artifact,
|
||
)
|
||
from deerflow.persistence.position_roundtable.delivery import (
|
||
session_intent_position_id as _chain_intent_position_id,
|
||
)
|
||
from deerflow.persistence.position_roundtable.sql import SessionWriteConflictError
|
||
from deerflow.runtime.user_context import get_effective_user_id
|
||
|
||
# The application mounts this router below both `/api` (canonical) and the
|
||
# long-established `/api/multi-agent` namespace (intranet compatibility).
|
||
router = APIRouter(prefix="/position-roundtable", tags=["position-roundtable"])
|
||
|
||
NODE_STATUSES = {"locked", "ready", "running", "done", "stale", "rejected", "error"}
|
||
_HANDOFF_TEXT_BUDGET = 12_000
|
||
_MAX_INLINE_ARTIFACTS_PER_NODE = 4
|
||
_MAX_INLINE_ARTIFACT_CHARS = 1_800
|
||
# Upstream deliverables are additionally mirrored into the current seat's own
|
||
# workspace (outputs/<rel> → workspace/upstream/<rel>) before each turn. Every
|
||
# seat runs on an isolated LangGraph thread, so without this copy a downstream
|
||
# ``read_file`` on the upstream artifact path fails with "File not found" and
|
||
# the seat is limited to the short excerpt inlined into the handoff prompt.
|
||
# The mirror deliberately stays outside outputs/ so it never appears in this
|
||
# seat's own delivery manifest.
|
||
_UPSTREAM_MIRROR_PREFIX = "mnt/user-data/workspace/upstream"
|
||
_MAX_MIRRORED_ARTIFACTS = 8
|
||
_MAX_MIRRORED_ARTIFACT_BYTES = 2_000_000
|
||
_MIRRORED_TOTAL_BYTES = 8_000_000
|
||
# The handoff prompt remains intentionally small, but the recovery copy stored
|
||
# in SQL must retain a complete normal Markdown/JSON deliverable. The ceiling
|
||
# protects the database from an accidental multi-gigabyte log while covering
|
||
# substantially larger reports than this workspace normally produces.
|
||
_MAX_PERSISTED_ARTIFACT_CHARS = 4_000_000
|
||
_BUILTIN_MATERIAL_BUDGET = 24_000
|
||
SUMMARY_AGENT_ID = "roundtable-summary"
|
||
ACTION_PLAN_AGENT_ID = "position-action-planner"
|
||
_BUILTIN_MARKDOWN_OUTPUTS = {
|
||
"summary": "mnt/user-data/outputs/岗位会商方案总结报告.md",
|
||
"action_plan": "mnt/user-data/outputs/岗位会商行动规划.md",
|
||
}
|
||
_TEXT_ARTIFACT_SUFFIXES = {
|
||
".csv",
|
||
".json",
|
||
".log",
|
||
".md",
|
||
".markdown",
|
||
".rst",
|
||
".txt",
|
||
".xml",
|
||
".yaml",
|
||
".yml",
|
||
}
|
||
|
||
|
||
class NodeResponse(BaseModel):
|
||
id: str
|
||
session_id: str
|
||
node_key: str
|
||
stage_index: int
|
||
seat_index: int
|
||
agent_id: str
|
||
position_id: str | None = None
|
||
thread_id: str | None = None
|
||
status: str
|
||
latest_answer: str | None = None
|
||
artifact_manifest: list[dict[str, Any]] = Field(default_factory=list)
|
||
conversation_snapshot: dict[str, Any] | None = None
|
||
revision: int = 0
|
||
# 节点乐观锁版本号(任意字段变更自增;revision 仅计真实交付)。
|
||
version: int = 0
|
||
upstream_revision: dict[str, Any] | None = None
|
||
rejection_count: int = 0
|
||
last_rejection: dict[str, Any] | None = None
|
||
invalidated_by: str | None = None
|
||
created_at: str | None = None
|
||
updated_at: str | None = None
|
||
|
||
|
||
class SessionResponse(BaseModel):
|
||
id: str
|
||
# Older personal rows can already contain a mock task id; task_scoped is
|
||
# the authoritative marker for task-bound rows. History is partitioned by
|
||
# user, so it only controls snapshot-id locking and list bucketing now.
|
||
external_task_id: str | None = None
|
||
task_scoped: bool = False
|
||
task_snapshot: dict[str, Any] = Field(default_factory=dict)
|
||
chain_id: str | None = None
|
||
chain_snapshot: dict[str, Any] | None = None
|
||
intent_thread_id: str | None = None
|
||
intent_snapshot: dict[str, Any] | None = None
|
||
intent_conversation: dict[str, Any] | None = None
|
||
summary_thread_id: str | None = None
|
||
summary_snapshot: dict[str, Any] | None = None
|
||
summary_conversation: dict[str, Any] | None = None
|
||
action_plan_thread_id: str | None = None
|
||
action_plan_snapshot: dict[str, Any] | None = None
|
||
action_plan_conversation: dict[str, Any] | None = None
|
||
status: str
|
||
# Session 级乐观锁版本号:每次命令式写操作自增,前端可回传做 CAS。
|
||
version: int = 0
|
||
created_at: str | None = None
|
||
updated_at: str | None = None
|
||
|
||
|
||
class SessionDetailResponse(SessionResponse):
|
||
nodes: list[NodeResponse] = Field(default_factory=list)
|
||
|
||
|
||
class SessionCreateRequest(BaseModel):
|
||
external_task_id: str | None = Field(default=None, max_length=128)
|
||
task_snapshot: dict[str, Any] = Field(default_factory=dict)
|
||
|
||
@field_validator("external_task_id")
|
||
@classmethod
|
||
def _normalize_external_task_id(cls, value: str | None) -> str | None:
|
||
if value is None:
|
||
return None
|
||
normalized = value.strip()
|
||
return normalized or None
|
||
|
||
|
||
class TaskSnapshotUpdateRequest(BaseModel):
|
||
task_snapshot: dict[str, Any] = Field(default_factory=dict)
|
||
|
||
|
||
class IntentUpdateRequest(BaseModel):
|
||
intent_thread_id: str | None = Field(default=None, max_length=128)
|
||
intent_snapshot: dict[str, Any] | None = None
|
||
conversation_snapshot: dict[str, Any] | None = None
|
||
|
||
|
||
class ConversationStateUpdateRequest(BaseModel):
|
||
"""Durable UI recovery state written while a normal chat is in progress."""
|
||
|
||
thread_id: str | None = Field(default=None, max_length=128)
|
||
conversation_snapshot: dict[str, Any] = Field(default_factory=dict)
|
||
# Client-generated idempotency key; one click/retry reuses the same id so
|
||
# the command executes exactly once. Absent → server generates one (legacy).
|
||
command_id: str | None = Field(default=None, max_length=64)
|
||
|
||
|
||
class ActivateRequest(BaseModel):
|
||
chain_id: str = Field(min_length=1, max_length=64)
|
||
command_id: str | None = Field(default=None, max_length=64)
|
||
|
||
|
||
class NodeThreadRequest(BaseModel):
|
||
thread_id: str = Field(min_length=1, max_length=128)
|
||
command_id: str | None = Field(default=None, max_length=64)
|
||
|
||
|
||
class CompleteTurnRequest(BaseModel):
|
||
latest_answer: str = Field(default="", max_length=500_000)
|
||
artifact_manifest: list[dict[str, Any]] = Field(default_factory=list)
|
||
upstream_revision: dict[str, Any] | None = None
|
||
command_id: str | None = Field(default=None, max_length=64)
|
||
|
||
|
||
class RejectNodeRequest(BaseModel):
|
||
# The UI only asks for an optional rejection reason. Keep ``title`` for
|
||
# compatible clients and for the audit record, but provide a stable
|
||
# default so neither field turns a valid rejection into a 422/400.
|
||
title: str = Field(default="产物驳回", max_length=120)
|
||
reason: str = Field(default="", max_length=5_000)
|
||
command_id: str | None = Field(default=None, max_length=64)
|
||
|
||
|
||
class ActionPlanDictionaryOption(BaseModel):
|
||
"""One task-system dictionary row, preserving dictValue as the model input."""
|
||
|
||
label: str = Field(default="", max_length=256)
|
||
value: str | int | float = ""
|
||
id: str | int | None = None
|
||
dict_type: str = Field(default="", max_length=128)
|
||
is_default: bool = False
|
||
|
||
|
||
class ActionPlanDictionaries(BaseModel):
|
||
"""The four action fields that must be selected from task-system dictionaries."""
|
||
|
||
publish_channels: list[ActionPlanDictionaryOption] = Field(default_factory=list)
|
||
action_types: list[ActionPlanDictionaryOption] = Field(default_factory=list)
|
||
info_formats: list[ActionPlanDictionaryOption] = Field(default_factory=list)
|
||
responsible_positions: list[ActionPlanDictionaryOption] = Field(default_factory=list)
|
||
|
||
|
||
class PrepareBuiltinRunRequest(BaseModel):
|
||
"""Optional follow-up instruction for one persisted built-in agent."""
|
||
|
||
message: str = Field(default="", max_length=20_000)
|
||
# 岗位会商会话允许本地演示使用非数字 taskId。真实任务系统新增行动
|
||
# 时由专属 Skill 再按外部接口契约校验正整数,不能在 prepare 阶段阻断报告生成。
|
||
task_id: str | int | None = Field(default=None, max_length=128)
|
||
action_dictionaries: ActionPlanDictionaries | None = None
|
||
|
||
|
||
class BuiltinRunPreparationResponse(BaseModel):
|
||
agent_id: str
|
||
display_name: str
|
||
thread_id: str | None = None
|
||
prompt: str
|
||
source_revisions: dict[str, int] = Field(default_factory=dict)
|
||
|
||
|
||
class CompleteBuiltinRunRequest(BaseModel):
|
||
thread_id: str = Field(min_length=1, max_length=128)
|
||
content: str = Field(default="", max_length=500_000)
|
||
artifact_manifest: list[dict[str, Any]] = Field(default_factory=list)
|
||
source_revisions: dict[str, int] = Field(default_factory=dict)
|
||
model: str | None = Field(default=None, max_length=256)
|
||
|
||
|
||
class DemoSeedRequest(BaseModel):
|
||
"""Canonical frontend fixture imported into the normal SQL session model."""
|
||
|
||
demo_key: str = Field(min_length=1, max_length=96)
|
||
session: dict[str, Any] = Field(default_factory=dict)
|
||
nodes: list[dict[str, Any]] = Field(default_factory=list, max_length=32)
|
||
|
||
|
||
def _current_user_id(request: Request) -> str:
|
||
user = getattr(request.state, "user", None)
|
||
return str(user.id) if user is not None else get_effective_user_id()
|
||
|
||
|
||
def _task_snapshot_for_scope(
|
||
external_task_id: str | None,
|
||
task_snapshot: dict[str, Any],
|
||
*,
|
||
task_scoped: bool = True,
|
||
) -> dict[str, Any]:
|
||
"""Lock the snapshot id only for an explicit task-bound scope."""
|
||
if not task_scoped or not external_task_id:
|
||
return dict(task_snapshot)
|
||
return {**task_snapshot, "id": external_task_id, "externalTaskId": external_task_id}
|
||
|
||
|
||
def _normalize_external_task_id(value: str | None) -> str | None:
|
||
if value is None:
|
||
return None
|
||
normalized = value.strip()
|
||
return normalized or None
|
||
|
||
|
||
def _get_store(request: Request):
|
||
store = getattr(request.app.state, "position_roundtable_store", None)
|
||
if store is None:
|
||
raise HTTPException(status_code=503, detail="Position roundtable store not available")
|
||
return store
|
||
|
||
|
||
def _get_chain_store(request: Request):
|
||
store = getattr(request.app.state, "roundtable_chain_store", None)
|
||
if store is None:
|
||
raise HTTPException(status_code=503, detail="Business chain store not available")
|
||
return store
|
||
|
||
|
||
async def _require_writable_session(
|
||
request: Request, session_id: str
|
||
) -> tuple[Any, str, dict[str, Any]]:
|
||
"""Load a Session that can still accept task, intent, or node updates."""
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
session = await store.get_session(session_id, user_id)
|
||
if session is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable session not found")
|
||
if session.get("status") == "archived":
|
||
raise HTTPException(status_code=409, detail="该岗位协同会话已归档,仅可查看历史记录")
|
||
return store, user_id, session
|
||
|
||
|
||
def _command_id_or_new(command_id: str | None) -> str:
|
||
"""Use the client's idempotency key, or mint one for legacy clients."""
|
||
value = (command_id or "").strip()
|
||
return value or uuid4().hex
|
||
|
||
|
||
def _command_result_or_raise(result: dict[str, Any]) -> dict[str, Any]:
|
||
"""Return an ok command result, or translate a rejected command into HTTP."""
|
||
if result.get("ok"):
|
||
return result
|
||
raise HTTPException(
|
||
status_code=int(result.get("http_status") or 409),
|
||
detail=result.get("detail") or "命令执行被拒绝",
|
||
)
|
||
|
||
|
||
async def _update_session_or_raise(
|
||
store: Any,
|
||
session_id: str,
|
||
user_id: str,
|
||
**changes: Any,
|
||
) -> dict[str, Any]:
|
||
"""Apply a locked session snapshot write and map stale-state races to 409."""
|
||
try:
|
||
row = await store.update_session(session_id, user_id, **changes)
|
||
except SessionWriteConflictError as exc:
|
||
raise HTTPException(status_code=409, detail=exc.detail) from exc
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable session not found")
|
||
return row
|
||
|
||
|
||
def _project_session(row: dict[str, Any]) -> SessionResponse:
|
||
return SessionResponse(**row)
|
||
|
||
|
||
def _project_node(row: dict[str, Any]) -> NodeResponse:
|
||
return NodeResponse(**row)
|
||
|
||
|
||
async def _load_detail(request: Request, session_id: str) -> SessionDetailResponse:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
row = await store.get_session(session_id, user_id)
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable session not found")
|
||
nodes = await store.list_nodes(session_id, user_id)
|
||
return SessionDetailResponse(**row, nodes=[_project_node(node) for node in nodes or []])
|
||
|
||
|
||
def _chain_nodes(
|
||
chain: dict[str, Any],
|
||
*,
|
||
available_agent_ids: set[str] | None = None,
|
||
) -> list[dict[str, Any]]:
|
||
seats = chain.get("seats") if isinstance(chain.get("seats"), list) else []
|
||
if not seats:
|
||
raise HTTPException(status_code=400, detail="业务链条没有可执行的智能体席位")
|
||
seat_by_agent_id = {str(seat.get("agent_id")): seat for seat in seats if seat.get("agent_id")}
|
||
if len(seat_by_agent_id) != len(seats):
|
||
raise HTTPException(status_code=400, detail="业务链条包含无效或重复的智能体席位")
|
||
|
||
raw_stages = chain.get("stages")
|
||
stages: list[list[str]]
|
||
if isinstance(raw_stages, list) and raw_stages:
|
||
stages = [[str(agent_id) for agent_id in stage] for stage in raw_stages if isinstance(stage, list)]
|
||
else:
|
||
stages = [[str(seat["agent_id"])] for seat in seats]
|
||
if not stages or any(not stage for stage in stages):
|
||
raise HTTPException(status_code=400, detail="业务链条阶段配置无效")
|
||
|
||
result: list[dict[str, Any]] = []
|
||
for stage_index, stage in enumerate(stages):
|
||
for seat_index, agent_id in enumerate(stage):
|
||
seat = seat_by_agent_id.get(agent_id)
|
||
if seat is None:
|
||
raise HTTPException(status_code=400, detail=f"业务链条阶段包含不存在的席位: {agent_id}")
|
||
if available_agent_ids is not None and agent_id not in available_agent_ids:
|
||
raise HTTPException(
|
||
status_code=400,
|
||
detail=f"业务链条席位「{seat.get('name') or agent_id}」缺少可用智能体,请先重新配置业务链条",
|
||
)
|
||
position_id = seat.get("position_id")
|
||
if not isinstance(position_id, str) or not position_id.strip():
|
||
raise HTTPException(
|
||
status_code=400,
|
||
detail=f"智能体「{seat.get('name') or agent_id}」尚未配置岗位,请先从岗位协同入口配置",
|
||
)
|
||
result.append(
|
||
{
|
||
"node_key": f"stage-{stage_index + 1}-seat-{seat_index + 1}-{agent_id}",
|
||
"stage_index": stage_index,
|
||
"seat_index": seat_index,
|
||
"agent_id": agent_id,
|
||
"position_id": position_id.strip(),
|
||
"status": "ready" if stage_index == 0 else "locked",
|
||
}
|
||
)
|
||
return result
|
||
|
||
|
||
async def _available_agent_ids(request: Request, agent_ids: set[str]) -> set[str] | None:
|
||
"""Return the configured agent ids if the agent store is available.
|
||
|
||
Position collaboration freezes a business-chain snapshot before execution.
|
||
Validating seat agent ids here gives the user a concrete configuration
|
||
error instead of failing later when a node tries to create its chat thread.
|
||
"""
|
||
agent_store = getattr(request.app.state, "agent_store", None)
|
||
if agent_store is None:
|
||
return None
|
||
available: set[str] = set()
|
||
for agent_id in agent_ids:
|
||
if await agent_store.get_any(agent_id) is not None:
|
||
available.add(agent_id)
|
||
return available
|
||
|
||
|
||
def _truncate_text(value: Any, limit: int) -> str:
|
||
if limit <= 0:
|
||
return ""
|
||
text = value if isinstance(value, str) else str(value or "")
|
||
if len(text) <= limit:
|
||
return text
|
||
return f"{text[:limit]}\n…(已按交接上下文长度截断)"
|
||
|
||
|
||
def _text_from_keys(source: dict[str, Any], keys: list[str], fallback: str = "") -> str:
|
||
for key in keys:
|
||
value = source.get(key)
|
||
if isinstance(value, str) and value.strip():
|
||
return value.strip()
|
||
return fallback
|
||
|
||
|
||
def _seat_meta(session: dict[str, Any], node: dict[str, Any]) -> tuple[str, str]:
|
||
chain = session.get("chain_snapshot") if isinstance(session.get("chain_snapshot"), dict) else {}
|
||
seats = chain.get("seats") if isinstance(chain.get("seats"), list) else []
|
||
seat = next(
|
||
(
|
||
item
|
||
for item in seats
|
||
if isinstance(item, dict) and str(item.get("agent_id") or "") == str(node.get("agent_id") or "")
|
||
),
|
||
{},
|
||
)
|
||
agent_name = str(seat.get("name") or node.get("agent_id") or "未命名智能体")
|
||
position_id = str(node.get("position_id") or seat.get("position_id") or "未配置岗位")
|
||
return agent_name, position_id
|
||
|
||
|
||
def _delivered_nodes(nodes: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||
return [node for node in nodes if node.get("status") == "done"]
|
||
|
||
|
||
def _all_nodes_delivered(
|
||
nodes: list[dict[str, Any]],
|
||
*,
|
||
intent_position_id: str | None = None,
|
||
) -> bool:
|
||
return bool(nodes) and all(
|
||
node.get("status") == "done"
|
||
and (_is_intent_position_node(node, intent_position_id) or _node_has_markdown_delivery(node))
|
||
for node in nodes
|
||
)
|
||
|
||
|
||
def _is_intent_position_node(
|
||
node: dict[str, Any],
|
||
intent_position_id: str | None = None,
|
||
) -> bool:
|
||
"""Node-dict wrapper over the shared harness decision (see delivery.py)."""
|
||
return _is_intent_position_by_id(
|
||
str(node.get("position_id") or "").strip(), intent_position_id
|
||
)
|
||
|
||
|
||
def _session_intent_position_id(session: dict[str, Any]) -> str | None:
|
||
"""Session-dict wrapper over the shared harness decision (see delivery.py)."""
|
||
return _chain_intent_position_id(session.get("chain_snapshot"))
|
||
|
||
|
||
async def _configured_intent_position_id(request: Request) -> str:
|
||
"""Get the one enabled special role for a newly activated session.
|
||
|
||
The chosen id is frozen into ``chain_snapshot`` on activation. A later
|
||
configuration change therefore applies to future workspaces without
|
||
reinterpreting completion rules in an in-flight historical business chain.
|
||
"""
|
||
store = getattr(request.app.state, "position_role_store", None)
|
||
if store is not None:
|
||
rows = await store.list_roles()
|
||
selected = next(
|
||
(
|
||
row.get("id")
|
||
for row in rows
|
||
if row.get("enabled") and row.get("role_type") == "intent" and isinstance(row.get("id"), str)
|
||
),
|
||
None,
|
||
)
|
||
if isinstance(selected, str) and selected.strip():
|
||
return selected.strip()
|
||
return "intelligence"
|
||
|
||
|
||
def _node_has_markdown_delivery(node: dict[str, Any]) -> bool:
|
||
artifacts = node.get("artifact_manifest")
|
||
return isinstance(artifacts, list) and any(_is_markdown_artifact(item) for item in artifacts)
|
||
|
||
|
||
def _intent_string_list(intent: dict[str, Any], key: str) -> list[str]:
|
||
"""Return a clean, display-safe string list from a persisted intent field."""
|
||
value = intent.get(key)
|
||
if not isinstance(value, list):
|
||
return []
|
||
return [item.strip() for item in value if isinstance(item, str) and item.strip()]
|
||
|
||
|
||
def _position_intent_sections(intent: Any) -> dict[str, list[str]]:
|
||
"""Read both the dedicated position-intent schema and legacy snapshots.
|
||
|
||
New position sessions persist the four position-specific fields. Mapping
|
||
the legacy Step-1 fields keeps already-created sessions executable after
|
||
this rollout instead of stranding them at the activation gate.
|
||
"""
|
||
if not isinstance(intent, dict):
|
||
return {
|
||
"coreGoals": [],
|
||
"riskWarnings": [],
|
||
"keyPoints": [],
|
||
"strategicSignificance": [],
|
||
}
|
||
has_position_schema = any(
|
||
key in intent
|
||
for key in ("coreGoals", "riskWarnings", "keyPoints", "strategicSignificance")
|
||
)
|
||
if has_position_schema:
|
||
return {
|
||
"coreGoals": _intent_string_list(intent, "coreGoals"),
|
||
"riskWarnings": _intent_string_list(intent, "riskWarnings"),
|
||
"keyPoints": _intent_string_list(intent, "keyPoints"),
|
||
"strategicSignificance": _intent_string_list(intent, "strategicSignificance"),
|
||
}
|
||
|
||
objective = _text_from_keys(intent, ["objective", "goal"])
|
||
return {
|
||
"coreGoals": _intent_string_list(intent, "constraints"),
|
||
"riskWarnings": [],
|
||
"keyPoints": _intent_string_list(intent, "assumptions"),
|
||
"strategicSignificance": [objective] if objective else [],
|
||
}
|
||
|
||
|
||
def _intent_node_answer(intent: Any) -> str:
|
||
if not isinstance(intent, dict):
|
||
return "情报分析岗已完成任务意图识别。"
|
||
sections = _position_intent_sections(intent)
|
||
lines = ["情报分析岗已完成任务意图识别。"]
|
||
labels = (
|
||
("核心目标", sections["coreGoals"]),
|
||
("风险提示", sections["riskWarnings"]),
|
||
("关键要点", sections["keyPoints"]),
|
||
("战略意义", sections["strategicSignificance"]),
|
||
)
|
||
for label, values in labels:
|
||
if values:
|
||
lines.append(f"{label}:")
|
||
lines.extend(f"- {item}" for item in values)
|
||
return "\n".join(lines)
|
||
|
||
|
||
def _has_confirmed_intent(intent: Any) -> bool:
|
||
"""Only structured Step-1 output may start a frozen business chain."""
|
||
if not isinstance(intent, dict):
|
||
return False
|
||
if any(
|
||
key in intent
|
||
for key in ("coreGoals", "riskWarnings", "keyPoints", "strategicSignificance")
|
||
):
|
||
required = ("coreGoals", "riskWarnings", "keyPoints", "strategicSignificance")
|
||
if not all(isinstance(intent.get(key), list) for key in required):
|
||
return False
|
||
sections = _position_intent_sections(intent)
|
||
return bool(sections["coreGoals"] or sections["strategicSignificance"])
|
||
objective = _text_from_keys(intent, ["objective", "goal"])
|
||
constraints = intent.get("constraints")
|
||
return bool(objective and isinstance(constraints, list))
|
||
|
||
|
||
def _prime_intent_nodes(
|
||
nodes: list[dict[str, Any]],
|
||
intent: Any,
|
||
*,
|
||
intent_position_id: str | None = None,
|
||
) -> list[dict[str, Any]]:
|
||
now_answer = _intent_node_answer(intent)
|
||
primed: list[dict[str, Any]] = []
|
||
for node in nodes:
|
||
next_node = dict(node)
|
||
if _is_intent_position_node(next_node, intent_position_id):
|
||
next_node.update(
|
||
{
|
||
"status": "done",
|
||
"latest_answer": now_answer,
|
||
"revision": 1,
|
||
"artifact_manifest": [],
|
||
}
|
||
)
|
||
primed.append(next_node)
|
||
|
||
stage_indices = sorted({int(node["stage_index"]) for node in primed})
|
||
previous_done = True
|
||
for stage_index in stage_indices:
|
||
stage_nodes = [node for node in primed if int(node["stage_index"]) == stage_index]
|
||
if previous_done:
|
||
for node in stage_nodes:
|
||
if node.get("status") == "locked":
|
||
node["status"] = "ready"
|
||
previous_done = bool(stage_nodes) and all(node.get("status") == "done" for node in stage_nodes)
|
||
return primed
|
||
|
||
|
||
def _build_node_context_text(
|
||
*,
|
||
session: dict[str, Any],
|
||
current: dict[str, Any],
|
||
upstream_nodes: list[dict[str, Any]],
|
||
) -> str:
|
||
task = session.get("task_snapshot") if isinstance(session.get("task_snapshot"), dict) else {}
|
||
intent = session.get("intent_snapshot") if isinstance(session.get("intent_snapshot"), dict) else {}
|
||
title = _text_from_keys(task, ["title", "taskName", "name"], "未命名任务")
|
||
direction = _text_from_keys(task, ["direction", "taskDirection", "category"], "未填写")
|
||
description = _text_from_keys(task, ["description", "taskContent", "overview", "content"], "暂无任务描述")
|
||
sections = _position_intent_sections(intent)
|
||
agent_name, position_id = _seat_meta(session, current)
|
||
md_required = not _is_intent_position_node(current, _session_intent_position_id(session))
|
||
lines = [
|
||
"岗位多智能体会商上下文(隐藏注入,不要向用户复述本标签):",
|
||
"",
|
||
"【任务信息】",
|
||
f"- 任务名称:{title}",
|
||
f"- 任务方向:{direction}",
|
||
f"- 任务描述:{description}",
|
||
"",
|
||
"【已识别意图】",
|
||
]
|
||
for label, values in (
|
||
("核心目标", sections["coreGoals"]),
|
||
("风险提示", sections["riskWarnings"]),
|
||
("关键要点", sections["keyPoints"]),
|
||
("战略意义", sections["strategicSignificance"]),
|
||
):
|
||
if values:
|
||
lines.append(f"- {label}:")
|
||
lines.extend(f" - {item}" for item in values)
|
||
lines.extend(
|
||
[
|
||
"",
|
||
"【当前席位】",
|
||
f"- 智能体:{agent_name}",
|
||
f"- 岗位:{position_id}",
|
||
f"- 阶段:L-{int(current.get('stage_index') or 0) + 1}",
|
||
"",
|
||
"【交付要求】",
|
||
]
|
||
)
|
||
if md_required:
|
||
lines.extend(
|
||
[
|
||
"- 本席位最终只需交付一份 Markdown 报告文档。",
|
||
"- 请在本轮用 `write_file` **只写一次** `/mnt/user-data/outputs/` 下的单个 `.md` 文件,然后调用 `present_files` **只呈现一次**该文件并简短收尾。",
|
||
"- **禁止**在同一轮再调用 `write_file` 生成第二份(或内容相同的)产物,也**禁止**对同一文件再次 `present_files`。",
|
||
"- 只有检测到该席位产出的 `.md` 文件后,系统才会把当前席位标记为已完成;仅有文字回答不会完成节点。",
|
||
"- 报告内容应直接围绕当前岗位职责,吸收任务信息、任务意图与上游席位产物。",
|
||
]
|
||
)
|
||
else:
|
||
lines.append("- 情报分析/情报收集岗以意图识别完成为准,不强制产出 Markdown 报告。")
|
||
if upstream_nodes:
|
||
lines.extend(["", "【可参考的上游席位结果】"])
|
||
for index, node in enumerate(upstream_nodes, start=1):
|
||
lines.append(
|
||
f"{index}. {node.get('agent_id')} / {node.get('position_id')} / revision={node.get('revision', 0)}"
|
||
)
|
||
answer = str(node.get("latest_answer") or "").strip()
|
||
if answer:
|
||
lines.append(f" 问答结果:{answer}")
|
||
artifacts = node.get("artifacts")
|
||
if isinstance(artifacts, list) and artifacts:
|
||
lines.append(" 产物:")
|
||
for artifact in artifacts[:5]:
|
||
if not isinstance(artifact, dict):
|
||
continue
|
||
artifact_name = artifact.get("name") or artifact.get("path") or "未命名产物"
|
||
lines.append(f" - {artifact_name}")
|
||
if artifact.get("mirrored_to_sandbox") and artifact.get("path"):
|
||
lines.append(
|
||
f" 完整文件已同步到本席位沙箱 /{artifact['path']},可直接用 read_file 读取全文"
|
||
"(仅供参考;本席位自己的交付物仍必须写入 /mnt/user-data/outputs/)。"
|
||
)
|
||
content = str(artifact.get("content") or "").strip()
|
||
if content:
|
||
lines.append(f" 摘要/正文片段:{content}")
|
||
else:
|
||
lines.extend(["", "【可参考的上游席位结果】", "- 当前节点暂无已完成上游席位。"])
|
||
return "\n".join(lines)
|
||
|
||
|
||
def _artifact_lines(node: dict[str, Any]) -> list[str]:
|
||
artifacts = node.get("artifact_manifest")
|
||
if not isinstance(artifacts, list) or not artifacts:
|
||
return []
|
||
lines: list[str] = []
|
||
for artifact in artifacts[:8]:
|
||
if not isinstance(artifact, dict):
|
||
continue
|
||
name = artifact.get("name") or artifact.get("filename") or artifact.get("path") or "未命名产物"
|
||
path = artifact.get("path") or artifact.get("artifact_url") or ""
|
||
lines.append(f" - {name}{f'({path})' if path else ''}")
|
||
return lines
|
||
|
||
|
||
def _artifact_refs(artifacts: Any) -> list[dict[str, Any]]:
|
||
if not isinstance(artifacts, list):
|
||
return []
|
||
refs: list[dict[str, Any]] = []
|
||
for artifact in artifacts[:20]:
|
||
if not isinstance(artifact, dict):
|
||
continue
|
||
ref = {
|
||
"name": artifact.get("name") or artifact.get("filename") or "未命名产物",
|
||
"path": artifact.get("path") or artifact.get("artifact_url") or "",
|
||
"mime_type": artifact.get("mimeType") or artifact.get("mime_type") or None,
|
||
}
|
||
# A database recovery payload is already constrained when it is first
|
||
# persisted. Preserve it here so downstream agents can still consume
|
||
# an upstream report after the owning sandbox/output disk is gone.
|
||
content = artifact.get("content")
|
||
if isinstance(content, str):
|
||
ref["content"] = content
|
||
if artifact.get("content_truncated"):
|
||
ref["content_truncated"] = True
|
||
refs.append(ref)
|
||
return refs
|
||
|
||
|
||
def _is_inline_text_artifact(artifact: dict[str, Any]) -> bool:
|
||
"""Return whether a listed output is safe and useful to pass as text."""
|
||
mime_type = artifact.get("mime_type")
|
||
if isinstance(mime_type, str) and mime_type.lower().startswith("text/"):
|
||
return True
|
||
path = artifact.get("path")
|
||
return isinstance(path, str) and Path(path).suffix.lower() in _TEXT_ARTIFACT_SUFFIXES
|
||
|
||
|
||
def _read_text_excerpt(path: Path, max_chars: int) -> tuple[str, bool] | None:
|
||
"""Read one bounded UTF-8 text artifact without loading the whole file."""
|
||
if max_chars <= 0:
|
||
return None
|
||
# Four bytes per UTF-8 code point is enough to obtain the requested
|
||
# character-sized excerpt without a potentially unbounded file read.
|
||
with path.open("rb") as file:
|
||
raw = file.read(max_chars * 4 + 1)
|
||
if b"\x00" in raw:
|
||
return None
|
||
try:
|
||
text = raw.decode("utf-8")
|
||
except UnicodeDecodeError as error:
|
||
if len(raw) == max_chars * 4 + 1 and error.start >= len(raw) - 4:
|
||
text = raw[: error.start].decode("utf-8")
|
||
else:
|
||
# Generated artifacts are expected to be UTF-8. Keep a non-UTF-8
|
||
# output as a reference only instead of corrupting agent context.
|
||
return None
|
||
truncated = len(raw) > max_chars * 4 or len(text) > max_chars
|
||
return _truncate_text(text, max_chars), truncated
|
||
|
||
|
||
async def _artifact_refs_with_content(
|
||
thread_id: str | None, artifacts: Any, text_budget: int
|
||
) -> tuple[list[dict[str, Any]], int]:
|
||
"""Attach bounded excerpts of text outputs, retaining binary files as refs.
|
||
|
||
Paths are always resolved against the upstream node's own thread. This is
|
||
deliberate: an artifact manifest is client-supplied metadata, so it must
|
||
never be allowed to select another thread's user-data directory.
|
||
"""
|
||
refs = _artifact_refs(artifacts)
|
||
if text_budget <= 0:
|
||
return refs, 0
|
||
|
||
consumed = 0
|
||
inline_count = 0
|
||
for ref in refs:
|
||
if inline_count >= _MAX_INLINE_ARTIFACTS_PER_NODE:
|
||
break
|
||
if not _is_inline_text_artifact(ref):
|
||
continue
|
||
virtual_path = ref.get("path")
|
||
if not isinstance(virtual_path, str):
|
||
continue
|
||
normalized_path = virtual_path.replace("\\", "/").lstrip("/")
|
||
if not normalized_path.startswith("mnt/user-data/outputs/"):
|
||
continue
|
||
available = text_budget - consumed
|
||
if available <= 0:
|
||
break
|
||
excerpt_limit = min(_MAX_INLINE_ARTIFACT_CHARS, available)
|
||
persisted_content = ref.get("content")
|
||
if isinstance(persisted_content, str):
|
||
content = _truncate_text(persisted_content, excerpt_limit)
|
||
ref["content"] = content
|
||
ref["content_truncated"] = bool(
|
||
ref.get("content_truncated") or len(persisted_content) > len(content)
|
||
)
|
||
consumed += len(content)
|
||
inline_count += 1
|
||
continue
|
||
if not thread_id:
|
||
continue
|
||
try:
|
||
actual_path = await aresolve_thread_virtual_path(thread_id, normalized_path)
|
||
if not await asyncio.to_thread(actual_path.is_file):
|
||
continue
|
||
excerpt = await asyncio.to_thread(_read_text_excerpt, actual_path, excerpt_limit)
|
||
except (HTTPException, OSError):
|
||
# The file may have been removed after the manifest was stored.
|
||
# Keep its metadata so the UI still shows the original output.
|
||
continue
|
||
if excerpt is None:
|
||
continue
|
||
content, truncated = excerpt
|
||
ref["content"] = content
|
||
ref["content_truncated"] = truncated
|
||
consumed += len(content)
|
||
inline_count += 1
|
||
return refs, consumed
|
||
|
||
|
||
def _outputs_relative_path(value: Any) -> str:
|
||
"""Return the ``outputs/``-relative posix path of a manifest artifact.
|
||
|
||
Empty string when the path is not a plain file directly under
|
||
``mnt/user-data/outputs/`` (wrong prefix, traversal segment, ...).
|
||
"""
|
||
path = str(value or "").replace("\\", "/").lstrip("/")
|
||
prefix = "mnt/user-data/outputs/"
|
||
if not path.startswith(prefix):
|
||
return ""
|
||
relative = path[len(prefix) :]
|
||
if not relative:
|
||
return ""
|
||
segments = relative.split("/")
|
||
if any(segment in ("", ".", "..") for segment in segments):
|
||
return ""
|
||
return relative
|
||
|
||
|
||
def _mirror_virtual_path(relative: str) -> str:
|
||
return f"{_UPSTREAM_MIRROR_PREFIX}/{relative}"
|
||
|
||
|
||
async def _mirror_upstream_artifacts(
|
||
current: dict[str, Any],
|
||
upstream_candidates: list[dict[str, Any]],
|
||
) -> dict[str, str]:
|
||
"""Copy upstream deliverables into the current seat's workspace sandbox.
|
||
|
||
Bytes are preferred from the upstream thread's own sandbox file; the SQL
|
||
recovery copy persisted in the manifest is the fallback when that disk is
|
||
gone. Only server-recorded upstream threads and manifest paths that pass
|
||
:func:`_outputs_relative_path` are read, and the target resolves through
|
||
the same thread virtual-path resolver, so a client-supplied manifest can
|
||
never point the copy outside a thread's user-data directory. Idempotent
|
||
per turn: an identically sized mirror is left untouched. Best effort —
|
||
failures simply keep the excerpt-only handoff behaviour.
|
||
"""
|
||
target_thread_id = str(current.get("thread_id") or "")
|
||
if not target_thread_id or not upstream_candidates:
|
||
return {}
|
||
|
||
mirrored: dict[str, str] = {}
|
||
total_bytes = 0
|
||
for node in upstream_candidates:
|
||
source_thread_id = str(node.get("thread_id") or "")
|
||
for artifact in node.get("artifact_manifest") or []:
|
||
if len(mirrored) >= _MAX_MIRRORED_ARTIFACTS or total_bytes >= _MIRRORED_TOTAL_BYTES:
|
||
return mirrored
|
||
if not isinstance(artifact, dict):
|
||
continue
|
||
original = str(artifact.get("path") or "")
|
||
relative = _outputs_relative_path(original)
|
||
if not relative:
|
||
continue
|
||
try:
|
||
payload: bytes | None = None
|
||
if source_thread_id:
|
||
source_path = await aresolve_thread_virtual_path(source_thread_id, original)
|
||
if await asyncio.to_thread(source_path.is_file):
|
||
stat = await asyncio.to_thread(source_path.stat)
|
||
if stat.st_size > _MAX_MIRRORED_ARTIFACT_BYTES:
|
||
continue
|
||
payload = await asyncio.to_thread(source_path.read_bytes)
|
||
if payload is None:
|
||
content = artifact.get("content")
|
||
if not isinstance(content, str) or not content:
|
||
continue
|
||
payload = content.encode("utf-8")
|
||
if len(payload) > _MAX_MIRRORED_ARTIFACT_BYTES:
|
||
continue
|
||
mirror_virtual = _mirror_virtual_path(relative)
|
||
target_path = await aresolve_thread_virtual_path(target_thread_id, mirror_virtual)
|
||
if not (
|
||
await asyncio.to_thread(target_path.is_file)
|
||
and (await asyncio.to_thread(target_path.stat)).st_size == len(payload)
|
||
):
|
||
await asyncio.to_thread(target_path.parent.mkdir, parents=True, exist_ok=True)
|
||
await asyncio.to_thread(target_path.write_bytes, payload)
|
||
mirrored[original] = mirror_virtual
|
||
total_bytes += len(payload)
|
||
except (HTTPException, OSError, ValueError):
|
||
# The upstream file may have been removed, or a provider may be
|
||
# unavailable. Skip it; the inline excerpt still carries the
|
||
# essential content into the handoff.
|
||
continue
|
||
return mirrored
|
||
|
||
|
||
async def _handoff_context(request: Request, session_id: str, node_key: str) -> dict[str, Any]:
|
||
"""Construct bounded, deterministic hidden context for one direct-agent turn."""
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
session = await store.get_session(session_id, user_id)
|
||
nodes = await store.list_nodes(session_id, user_id)
|
||
if session is None or nodes is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable session not found")
|
||
if session.get("status") == "archived":
|
||
raise HTTPException(status_code=409, detail="该岗位协同会话已归档,仅可查看历史记录")
|
||
current = next((node for node in nodes if node["node_key"] == node_key), None)
|
||
if current is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable node not found")
|
||
if current["status"] == "locked":
|
||
raise HTTPException(status_code=409, detail="上游节点尚未完成,当前节点未解锁")
|
||
if current.get("invalidated_by"):
|
||
raise HTTPException(status_code=409, detail="上游产物已被驳回或重做,当前节点需等待前序节点重新完成")
|
||
|
||
# Same-stage nodes are intentionally excluded. The business chain only
|
||
# hands off completed previous stages, in stage/seat order.
|
||
remaining = _HANDOFF_TEXT_BUDGET
|
||
upstream_candidates = [
|
||
node
|
||
for node in nodes
|
||
if node["stage_index"] < current["stage_index"]
|
||
and node["status"] == "done"
|
||
]
|
||
upstream_nodes: list[dict[str, Any]] = []
|
||
for index, node in enumerate(upstream_candidates):
|
||
candidates_left = len(upstream_candidates) - index
|
||
per_node_limit = min(2_800, max(0, (remaining // candidates_left) * 2 // 3))
|
||
answer = _truncate_text(node.get("latest_answer"), per_node_limit)
|
||
remaining -= len(answer)
|
||
artifact_budget = min(
|
||
_MAX_INLINE_ARTIFACT_CHARS,
|
||
max(0, remaining // candidates_left),
|
||
)
|
||
artifacts, consumed = await _artifact_refs_with_content(
|
||
node.get("thread_id"), node.get("artifact_manifest"), artifact_budget
|
||
)
|
||
remaining -= consumed
|
||
upstream_nodes.append(
|
||
{
|
||
"node_key": node["node_key"],
|
||
"agent_id": node["agent_id"],
|
||
"position_id": node.get("position_id"),
|
||
"revision": node.get("revision", 0),
|
||
"latest_answer": answer,
|
||
"artifacts": artifacts,
|
||
}
|
||
)
|
||
if remaining <= 0:
|
||
break
|
||
|
||
# Point each artifact ref at the workspace mirror when one exists: the
|
||
# agent then reads a path that is actually resolvable inside its own
|
||
# thread sandbox instead of the upstream thread's outputs directory.
|
||
mirrored_paths = await _mirror_upstream_artifacts(current, upstream_candidates)
|
||
for node in upstream_nodes:
|
||
for artifact in node.get("artifacts") or []:
|
||
mirror = mirrored_paths.get(str(artifact.get("path") or ""))
|
||
if not mirror:
|
||
continue
|
||
artifact["upstream_path"] = artifact["path"]
|
||
artifact["path"] = mirror
|
||
artifact["mirrored_to_sandbox"] = True
|
||
|
||
return {
|
||
"kind": "position_roundtable_handoff",
|
||
"instructions": _build_node_context_text(
|
||
session=session,
|
||
current=current,
|
||
upstream_nodes=upstream_nodes,
|
||
),
|
||
"delivery_requirement": {
|
||
"type": "markdown_report",
|
||
"required": not _is_intent_position_node(current, _session_intent_position_id(session)),
|
||
"completion_rule": "普通业务席位必须在 /mnt/user-data/outputs/ 下交付一份 .md 报告文档才会标记完成;情报分析/情报收集岗以意图识别完成为准。",
|
||
},
|
||
"task": session.get("task_snapshot") or {},
|
||
"intent": session.get("intent_snapshot") or {},
|
||
"current_node": {
|
||
"node_key": current["node_key"],
|
||
"agent_id": current["agent_id"],
|
||
"position_id": current.get("position_id"),
|
||
"stage_index": current["stage_index"],
|
||
"seat_index": current["seat_index"],
|
||
},
|
||
"upstream_nodes": upstream_nodes,
|
||
"upstream_revision": {
|
||
node["node_key"]: node.get("revision", 0) for node in upstream_nodes
|
||
},
|
||
}
|
||
|
||
|
||
def _source_revisions(nodes: list[dict[str, Any]]) -> dict[str, int]:
|
||
"""Stable revision fence for a built-in report/planning run.
|
||
|
||
The report agents run asynchronously in the browser. If a user rejects or
|
||
re-delivers a node while a run is still streaming, this fence prevents that
|
||
old run from overwriting the newer business-chain state when it completes.
|
||
"""
|
||
return {str(node["node_key"]): int(node.get("revision") or 0) for node in _delivered_nodes(nodes)}
|
||
|
||
|
||
def _artifact_manifest_refs(artifacts: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||
"""Persist display metadata plus any database recovery payload."""
|
||
refs: list[dict[str, Any]] = []
|
||
for artifact in artifacts:
|
||
if not isinstance(artifact, dict):
|
||
continue
|
||
path = str(artifact.get("path") or artifact.get("artifact_url") or "").strip()
|
||
if not path.replace("\\", "/").lstrip("/").startswith("mnt/user-data/outputs/"):
|
||
continue
|
||
ref = {
|
||
"name": str(artifact.get("name") or artifact.get("filename") or path.rsplit("/", 1)[-1]),
|
||
"path": path,
|
||
"mime_type": artifact.get("mime_type") or artifact.get("mimeType") or None,
|
||
"modified_at": artifact.get("modified_at") or artifact.get("modifiedAt") or None,
|
||
"size_bytes": artifact.get("size_bytes") or artifact.get("sizeBytes") or None,
|
||
}
|
||
content = artifact.get("content")
|
||
if isinstance(content, str):
|
||
ref["content"] = content
|
||
if artifact.get("content_truncated"):
|
||
ref["content_truncated"] = True
|
||
refs.append(ref)
|
||
return refs
|
||
|
||
|
||
async def _persist_artifact_contents(
|
||
thread_id: str | None,
|
||
artifacts: list[dict[str, Any]],
|
||
) -> list[dict[str, Any]]:
|
||
"""Copy text deliverables into SQL so history survives ephemeral disks."""
|
||
refs = _artifact_manifest_refs(artifacts)
|
||
if not thread_id:
|
||
return refs
|
||
|
||
enriched: list[dict[str, Any]] = []
|
||
for artifact in refs:
|
||
path = str(artifact.get("path") or "")
|
||
suffix = Path(path).suffix.lower()
|
||
if suffix not in _TEXT_ARTIFACT_SUFFIXES:
|
||
enriched.append(artifact)
|
||
continue
|
||
try:
|
||
resolved = await aresolve_thread_virtual_path(thread_id, path)
|
||
content = await asyncio.to_thread(resolved.read_text, encoding="utf-8", errors="replace")
|
||
except Exception:
|
||
# Artifact enrichment is recovery hardening, not a reason to fail
|
||
# an otherwise successful delivery. Keep the manifest metadata if
|
||
# the sandbox disappeared or its resolver/provider is unavailable.
|
||
enriched.append(artifact)
|
||
continue
|
||
if len(content) > _MAX_PERSISTED_ARTIFACT_CHARS:
|
||
artifact = {
|
||
**artifact,
|
||
"content": content[:_MAX_PERSISTED_ARTIFACT_CHARS],
|
||
"content_truncated": True,
|
||
}
|
||
else:
|
||
artifact = {**artifact, "content": content}
|
||
enriched.append(artifact)
|
||
return enriched
|
||
|
||
|
||
async def _write_builtin_markdown_fallback(
|
||
*, thread_id: str, kind: str, content: str
|
||
) -> dict[str, Any]:
|
||
"""Materialize a built-in agent's visible answer when it skipped ``write_file``.
|
||
|
||
Position roundtable runs are direct chat streams. A model can therefore
|
||
finish with a complete visible report but still omit its required artifact
|
||
tool call. Keep the chat answer as the source of truth and mirror it into
|
||
the owning thread's outputs directory so the report remains previewable and
|
||
downloadable after the stream has ended.
|
||
"""
|
||
virtual_path = _BUILTIN_MARKDOWN_OUTPUTS.get(kind)
|
||
report = content.strip()
|
||
if virtual_path is None:
|
||
raise HTTPException(status_code=400, detail=f"Unknown built-in delivery kind: {kind}")
|
||
if not report:
|
||
raise HTTPException(status_code=422, detail="内置智能体未返回可保存为报告的正文")
|
||
|
||
try:
|
||
output_path = await aresolve_thread_virtual_path(thread_id, virtual_path)
|
||
await asyncio.to_thread(output_path.parent.mkdir, parents=True, exist_ok=True)
|
||
await asyncio.to_thread(
|
||
output_path.write_text,
|
||
report if report.endswith("\n") else f"{report}\n",
|
||
encoding="utf-8",
|
||
)
|
||
stat = await asyncio.to_thread(output_path.stat)
|
||
except OSError as error:
|
||
raise HTTPException(status_code=500, detail="无法保存内置智能体 Markdown 报告") from error
|
||
|
||
return {
|
||
"name": output_path.name,
|
||
"path": virtual_path,
|
||
"mime_type": "text/markdown",
|
||
"modified_at": datetime.fromtimestamp(stat.st_mtime, tz=UTC).isoformat(),
|
||
"size_bytes": stat.st_size,
|
||
}
|
||
|
||
|
||
async def _effective_delivery_material(
|
||
session: dict[str, Any], nodes: list[dict[str, Any]]
|
||
) -> str:
|
||
"""Build bounded source material from valid current deliveries only.
|
||
|
||
Stale/rejected rows are deliberately excluded. The persisted node answer
|
||
provides a quick synopsis; the Markdown files provide the actual
|
||
deliverable body, preferably from its database recovery copy and otherwise
|
||
through the owning thread's virtual-path resolver.
|
||
"""
|
||
remaining = _BUILTIN_MATERIAL_BUDGET
|
||
lines: list[str] = ["【岗位有效交付】"]
|
||
for index, node in enumerate(_delivered_nodes(nodes), start=1):
|
||
agent_name, position_id = _seat_meta(session, node)
|
||
lines.extend(
|
||
[
|
||
"",
|
||
f"### {index}. {agent_name}",
|
||
f"- 岗位:{position_id}",
|
||
f"- 节点:{node.get('node_key')}",
|
||
f"- 有效版本:v{node.get('revision', 0)}",
|
||
]
|
||
)
|
||
answer_limit = min(2_000, max(0, remaining // max(1, len(nodes))))
|
||
answer = _truncate_text(node.get("latest_answer"), answer_limit)
|
||
if answer:
|
||
remaining -= len(answer)
|
||
lines.extend(["- 本轮问答结论:", answer])
|
||
|
||
artifacts, consumed = await _artifact_refs_with_content(
|
||
node.get("thread_id"), node.get("artifact_manifest"), max(0, remaining)
|
||
)
|
||
remaining -= consumed
|
||
markdown_artifacts = [artifact for artifact in artifacts if _is_markdown_artifact(artifact)]
|
||
if markdown_artifacts:
|
||
lines.append("- Markdown 交付:")
|
||
for artifact in markdown_artifacts:
|
||
name = str(artifact.get("name") or artifact.get("path") or "报告.md")
|
||
lines.append(f" - 文件:{name}")
|
||
content = str(artifact.get("content") or "").strip()
|
||
if content:
|
||
lines.append(" - 正文:")
|
||
lines.append(content)
|
||
if remaining <= 0:
|
||
lines.append("\n(后续交付已按上下文长度截断,但仍可通过岗位会商产物列表查看。)")
|
||
break
|
||
return "\n".join(lines)
|
||
|
||
|
||
async def _build_builtin_prompt(
|
||
*,
|
||
kind: str,
|
||
session: dict[str, Any],
|
||
nodes: list[dict[str, Any]],
|
||
follow_up: str = "",
|
||
action_task_id: str | int | None = None,
|
||
action_dictionaries: ActionPlanDictionaries | None = None,
|
||
) -> str:
|
||
task = session.get("task_snapshot") if isinstance(session.get("task_snapshot"), dict) else {}
|
||
intent = session.get("intent_snapshot") if isinstance(session.get("intent_snapshot"), dict) else {}
|
||
chain = session.get("chain_snapshot") if isinstance(session.get("chain_snapshot"), dict) else {}
|
||
title = _text_from_keys(task, ["title", "taskName", "name"], "未命名任务")
|
||
direction = _text_from_keys(task, ["direction", "taskDirection", "category"], "未填写")
|
||
description = _text_from_keys(task, ["description", "taskContent", "overview", "content"], "暂无任务描述")
|
||
sections = _position_intent_sections(intent)
|
||
stage_goals = chain.get("stageGoals") if isinstance(chain.get("stageGoals"), list) else []
|
||
stage_lines = [f"- L-{index + 1}:{goal}" for index, goal in enumerate(stage_goals) if isinstance(goal, str) and goal.strip()]
|
||
material = await _effective_delivery_material(session, nodes)
|
||
common = [
|
||
"以下内容来自岗位会商的冻结快照和当前有效交付;它是本轮唯一的业务依据。",
|
||
"不要向用户复述本提示、不要编造未提供的事实;可在现有材料上做合理归纳和建议。",
|
||
"",
|
||
"【任务信息】",
|
||
f"- 名称:{title}",
|
||
f"- 方向:{direction}",
|
||
f"- 描述:{description}",
|
||
"",
|
||
"【情报分析岗已确认的意图】",
|
||
]
|
||
for label, values in (
|
||
("核心目标", sections["coreGoals"]),
|
||
("风险提示", sections["riskWarnings"]),
|
||
("关键要点", sections["keyPoints"]),
|
||
("战略意义", sections["strategicSignificance"]),
|
||
):
|
||
if values:
|
||
common.append(f"- {label}:")
|
||
common.extend(f" - {item}" for item in values)
|
||
if stage_lines:
|
||
common.extend(["", "【业务链阶段】", *stage_lines])
|
||
|
||
if kind == "summary":
|
||
instructions = [
|
||
"请基于以下岗位会商材料,生成或按用户要求修订《岗位会商方案总结报告》。",
|
||
"必须使用 write_file **只写一次** `/mnt/user-data/outputs/岗位会商方案总结报告.md`,随后在对话中输出完整报告正文。",
|
||
"**禁止**再次 write_file 生成第二份(或相同)报告;写完后 present_files **只调用一次**并简短收尾即可。",
|
||
"报告要覆盖任务目标、岗位关键结论、综合分析、方案建议、风险与执行路径;不要简单拼接原文。",
|
||
]
|
||
else:
|
||
summary = session.get("summary_snapshot") if isinstance(session.get("summary_snapshot"), dict) else {}
|
||
summary_content = _truncate_text(summary.get("content"), 12_000)
|
||
dictionary_groups = (
|
||
("发布渠道(publishChannel)", action_dictionaries.publish_channels if action_dictionaries else []),
|
||
("行动类型(actionType)", action_dictionaries.action_types if action_dictionaries else []),
|
||
("信息形式(infoFormat)", action_dictionaries.info_formats if action_dictionaries else []),
|
||
("责任岗位(responsiblePosition)", action_dictionaries.responsible_positions if action_dictionaries else []),
|
||
)
|
||
dictionaries_ready = bool(action_task_id is not None and all(options for _, options in dictionary_groups))
|
||
dictionary_lines = [
|
||
"【前端传入的行动字典(名称 → 必须原样使用的 dictValue)】",
|
||
f"- 当前 taskId:{action_task_id if action_task_id is not None else '(本地预览)'}",
|
||
]
|
||
for label, options in dictionary_groups:
|
||
dictionary_lines.append(f"- {label}:")
|
||
dictionary_lines.extend(
|
||
f" - {option.label} → {option.value}"
|
||
for option in options
|
||
)
|
||
dictionary_instruction = (
|
||
"actionType、publishChannel、infoFormat、responsiblePosition 必须从下方同名字典复制 dictValue 并写成整数;"
|
||
"taskId 必须原样写为当前 taskId。不要保留 id、title、goal、basis、owner_position、dependencies 等旧字段,"
|
||
"也不要用 0 或空值代替未选择的字典项。"
|
||
if dictionaries_ready
|
||
else "当前四类任务系统字典未获取完整,这是本地预览/接口未就绪状态。仍需生成 9 字段 JSON,"
|
||
"但四个枚举字段可暂填 0 占位,且不得调用 position-roundtable-action-create Skill 或外部任务系统。"
|
||
)
|
||
submit_instruction = (
|
||
"若当前 taskId 是正整数,首次生成完整行动清单后按行动规划智能体的 SOUL 调用 position-roundtable-action-create Skill 批量新增行动;"
|
||
"若 taskId 为本地演示标识或为空,只生成报告和 JSON,明确说明未调用任务系统。Skill 返回失败时如实报告,不能伪称已入库。"
|
||
if dictionaries_ready
|
||
else "本轮字典不完整,只生成报告和 JSON 预览;不得调用 Skill、不得请求外部任务系统,也不能宣称行动已入库。"
|
||
)
|
||
instructions = [
|
||
"请基于以下方案总结和岗位会商材料,生成或按用户要求修订《岗位会商行动规划》。",
|
||
"必须使用 write_file 分别写入两个最终交付:`/mnt/user-data/outputs/行动规划报告.md` 与 `/mnt/user-data/outputs/action-plan-subtasks.json`。",
|
||
"Markdown 报告需给出分阶段推进建议、优先级、风险控制与验收安排。"
|
||
"JSON 必须是严格合法数组;每个对象只能包含 actionName、actionDescription、actionType、"
|
||
"publishChannel、infoFormat、responsiblePosition、startTime、endTime、taskId 这 9 个字段。",
|
||
dictionary_instruction,
|
||
submit_instruction,
|
||
"",
|
||
"【已生效方案总结】",
|
||
summary_content or "(方案总结正文未写入快照,请依据下方有效岗位交付重新归纳。)",
|
||
"",
|
||
*dictionary_lines,
|
||
]
|
||
if follow_up.strip():
|
||
instructions.extend(["", "【用户本次补充/修改要求】", follow_up.strip()])
|
||
return "\n".join([*instructions, "", *common, "", material])
|
||
|
||
|
||
def _build_builtin_snapshot(
|
||
*,
|
||
kind: str,
|
||
thread_id: str,
|
||
content: str,
|
||
artifacts: list[dict[str, Any]],
|
||
source_revisions: dict[str, int],
|
||
model: str | None,
|
||
previous: dict[str, Any] | None = None,
|
||
) -> dict[str, Any]:
|
||
markdown = next((artifact for artifact in artifacts if _is_markdown_artifact(artifact)), None)
|
||
json_artifact = next((artifact for artifact in artifacts if _is_json_artifact(artifact)), None)
|
||
primary = markdown if kind == "summary" else markdown
|
||
previous_primary = previous.get("primary_artifact") if isinstance(previous, dict) else None
|
||
primary_changed = not isinstance(previous_primary, dict) or not isinstance(primary, dict) or (
|
||
_artifact_signature(previous_primary) != _artifact_signature(primary)
|
||
)
|
||
previous_content = str(previous.get("content") or "") if isinstance(previous, dict) else ""
|
||
return {
|
||
"kind": kind,
|
||
"status": "done",
|
||
"title": "岗位会商方案总结报告" if kind == "summary" else "岗位会商行动规划",
|
||
# A pure follow-up may only answer a question and leave the report file
|
||
# unchanged. Keep the last report body in that case while retaining
|
||
# the answer separately for audit/display.
|
||
"content": content.strip() if primary_changed or not previous_content else previous_content,
|
||
"last_response": content.strip(),
|
||
"thread_id": thread_id,
|
||
"artifact_manifest": _artifact_manifest_refs(artifacts),
|
||
"primary_artifact": primary,
|
||
"subtasks_artifact": json_artifact if kind == "action_plan" else None,
|
||
"source_revisions": source_revisions,
|
||
"model": model or None,
|
||
"generated_at": datetime.now(UTC).isoformat(),
|
||
}
|
||
|
||
|
||
@router.post("/demo/seed", response_model=SessionDetailResponse)
|
||
async def seed_demo_session(
|
||
request: Request, body: DemoSeedRequest
|
||
) -> SessionDetailResponse:
|
||
"""Idempotently materialise the frontend demo JSON for the current user.
|
||
|
||
The endpoint intentionally accepts a complete recovery snapshot instead of
|
||
inventing a second demo-only schema. Imported rows therefore exercise the
|
||
same history, position switching, artifact and conversation code paths as
|
||
real collaboration sessions.
|
||
"""
|
||
|
||
row, nodes = await _get_store(request).ensure_demo_session(
|
||
_current_user_id(request),
|
||
demo_key=body.demo_key,
|
||
session_payload=body.session,
|
||
nodes=body.nodes,
|
||
)
|
||
return SessionDetailResponse(
|
||
**row,
|
||
nodes=[_project_node(node) for node in nodes],
|
||
)
|
||
|
||
|
||
@router.get("/sessions", response_model=list[SessionResponse])
|
||
async def list_sessions(
|
||
request: Request,
|
||
external_task_id: str | None = Query(default=None, max_length=128),
|
||
) -> list[SessionResponse]:
|
||
task_id = _normalize_external_task_id(external_task_id)
|
||
rows = await _get_store(request).list_sessions(
|
||
_current_user_id(request), external_task_id=task_id
|
||
)
|
||
return [_project_session(row) for row in rows]
|
||
|
||
|
||
@router.post("/sessions", response_model=SessionResponse, status_code=201)
|
||
async def create_session(request: Request, body: SessionCreateRequest) -> SessionResponse:
|
||
row = await _get_store(request).create_session(
|
||
_current_user_id(request),
|
||
external_task_id=body.external_task_id,
|
||
task_snapshot=_task_snapshot_for_scope(
|
||
body.external_task_id,
|
||
body.task_snapshot,
|
||
task_scoped=bool(body.external_task_id),
|
||
),
|
||
task_scoped=bool(body.external_task_id),
|
||
)
|
||
return _project_session(row)
|
||
|
||
|
||
@router.get("/sessions/{session_id}", response_model=SessionDetailResponse)
|
||
async def get_session(request: Request, session_id: str) -> SessionDetailResponse:
|
||
return await _load_detail(request, session_id)
|
||
|
||
|
||
@router.put("/sessions/{session_id}/task", response_model=SessionResponse)
|
||
async def update_task_snapshot(
|
||
request: Request, session_id: str, body: TaskSnapshotUpdateRequest
|
||
) -> SessionResponse:
|
||
store, user_id, current = await _require_writable_session(request, session_id)
|
||
row = await _update_session_or_raise(
|
||
store,
|
||
session_id,
|
||
user_id,
|
||
task_snapshot=_task_snapshot_for_scope(
|
||
current.get("external_task_id"),
|
||
body.task_snapshot,
|
||
task_scoped=bool(current.get("task_scoped")),
|
||
),
|
||
allowed_statuses={"intent_pending"},
|
||
)
|
||
return _project_session(row)
|
||
|
||
|
||
@router.put("/sessions/{session_id}/intent", response_model=SessionResponse)
|
||
async def update_intent_snapshot(
|
||
request: Request, session_id: str, body: IntentUpdateRequest
|
||
) -> SessionResponse:
|
||
store, user_id, _ = await _require_writable_session(request, session_id)
|
||
row = await _update_session_or_raise(
|
||
store,
|
||
session_id,
|
||
user_id,
|
||
intent_thread_id=body.intent_thread_id,
|
||
intent_snapshot=body.intent_snapshot,
|
||
intent_conversation=body.conversation_snapshot,
|
||
allowed_statuses={"intent_pending"},
|
||
)
|
||
return _project_session(row)
|
||
|
||
|
||
@router.put("/sessions/{session_id}/summary/state", response_model=SessionResponse)
|
||
async def update_summary_conversation(
|
||
request: Request, session_id: str, body: ConversationStateUpdateRequest
|
||
) -> SessionResponse:
|
||
store, user_id, _ = await _require_writable_session(request, session_id)
|
||
row = await _update_session_or_raise(
|
||
store,
|
||
session_id,
|
||
user_id,
|
||
summary_thread_id=body.thread_id,
|
||
summary_conversation=body.conversation_snapshot,
|
||
)
|
||
return _project_session(row)
|
||
|
||
|
||
@router.put("/sessions/{session_id}/action-plan/state", response_model=SessionResponse)
|
||
async def update_action_plan_conversation(
|
||
request: Request, session_id: str, body: ConversationStateUpdateRequest
|
||
) -> SessionResponse:
|
||
store, user_id, _ = await _require_writable_session(request, session_id)
|
||
row = await _update_session_or_raise(
|
||
store,
|
||
session_id,
|
||
user_id,
|
||
action_plan_thread_id=body.thread_id,
|
||
action_plan_conversation=body.conversation_snapshot,
|
||
)
|
||
return _project_session(row)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/activate", response_model=SessionDetailResponse)
|
||
async def activate_session(
|
||
request: Request, session_id: str, body: ActivateRequest
|
||
) -> SessionDetailResponse:
|
||
store, user_id, current = await _require_writable_session(request, session_id)
|
||
if not _has_confirmed_intent(current.get("intent_snapshot")):
|
||
raise HTTPException(status_code=400, detail="请先完成并保存任务意图")
|
||
|
||
chain = await _get_chain_store(request).get_chain(
|
||
body.chain_id,
|
||
user_id,
|
||
is_admin=getattr(getattr(request.state, "user", None), "system_role", None) == "admin",
|
||
)
|
||
if chain is None:
|
||
raise HTTPException(status_code=404, detail="业务链条不存在或无权访问")
|
||
seat_ids = {
|
||
str(seat.get("agent_id"))
|
||
for seat in (chain.get("seats") if isinstance(chain.get("seats"), list) else [])
|
||
if seat.get("agent_id")
|
||
}
|
||
# Freeze the currently configured special role with this session. This
|
||
# prevents a later global config edit from changing the completion rule of
|
||
# an already running historical chain.
|
||
chain_snapshot = dict(chain)
|
||
intent_position_id = await _configured_intent_position_id(request)
|
||
chain_snapshot["_intent_position_id"] = intent_position_id
|
||
nodes = _chain_nodes(
|
||
chain_snapshot,
|
||
available_agent_ids=await _available_agent_ids(request, seat_ids),
|
||
)
|
||
nodes = _prime_intent_nodes(
|
||
nodes,
|
||
current.get("intent_snapshot"),
|
||
intent_position_id=intent_position_id,
|
||
)
|
||
# 命令式状态机:锁 session + 锁内判定「已激活→链相同幂等/链不同 409」,
|
||
# 消除原 activate_session 的 check-then-insert 竞态(并发双击 → 唯一键冲突 500)。
|
||
result = _command_result_or_raise(
|
||
await store.apply_session_command(
|
||
session_id,
|
||
user_id,
|
||
command_id=_command_id_or_new(body.command_id),
|
||
command_type="activate",
|
||
payload={
|
||
"chain_id": body.chain_id,
|
||
"chain_snapshot": chain_snapshot,
|
||
"nodes": nodes,
|
||
},
|
||
)
|
||
)
|
||
return SessionDetailResponse(
|
||
**result["session"],
|
||
nodes=[_project_node(node) for node in result["nodes"]],
|
||
)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/archive", response_model=SessionResponse)
|
||
async def archive_session(request: Request, session_id: str) -> SessionResponse:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
# 命令式状态机:已归档幂等返回,锁内串行化并发归档。每次操作一个稳定的
|
||
# command_id(URL 无 body),用 session 级固定派生键,重复请求天然重放。
|
||
result = _command_result_or_raise(
|
||
await store.apply_session_command(
|
||
session_id,
|
||
user_id,
|
||
command_id=f"archive-{session_id}",
|
||
command_type="archive",
|
||
payload={},
|
||
)
|
||
)
|
||
return _project_session(result["session"])
|
||
|
||
|
||
@router.delete("/sessions/{session_id}")
|
||
async def delete_session(request: Request, session_id: str) -> dict[str, bool]:
|
||
store = _get_store(request)
|
||
user_id = _current_user_id(request)
|
||
# 命令式状态机:命令记录无外键(保留供重放),删除后同 command_id 重试 → 重放
|
||
# ``{deleted: True}``,避免已删会话再次请求报 404 让前端误判失败。
|
||
_command_result_or_raise(
|
||
await store.apply_session_command(
|
||
session_id,
|
||
user_id,
|
||
command_id=f"delete-{session_id}",
|
||
command_type="delete",
|
||
payload={},
|
||
)
|
||
)
|
||
return {"deleted": True}
|
||
|
||
|
||
@router.get("/sessions/{session_id}/nodes/{node_key}", response_model=NodeResponse)
|
||
async def get_node(request: Request, session_id: str, node_key: str) -> NodeResponse:
|
||
row = await _get_store(request).get_node(session_id, node_key, _current_user_id(request))
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable node not found")
|
||
return _project_node(row)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/nodes/{node_key}/thread", response_model=NodeResponse)
|
||
async def bind_node_thread(
|
||
request: Request, session_id: str, node_key: str, body: NodeThreadRequest
|
||
) -> NodeResponse:
|
||
store, user_id, _ = await _require_writable_session(request, session_id)
|
||
# 命令式状态机:锁内判定「是否解锁 / 是否已被其它线程占用」并绑定。
|
||
# 原实现在事务外读取节点状态做判定(TOCTOU),并发下两个请求可同时通过校验。
|
||
result = _command_result_or_raise(
|
||
await store.apply_node_command(
|
||
session_id,
|
||
node_key,
|
||
user_id,
|
||
command_id=_command_id_or_new(body.command_id),
|
||
command_type="bind",
|
||
payload={"thread_id": body.thread_id},
|
||
)
|
||
)
|
||
return _project_node(result["node"])
|
||
|
||
|
||
@router.put(
|
||
"/sessions/{session_id}/nodes/{node_key}/state",
|
||
response_model=NodeResponse,
|
||
)
|
||
async def update_node_conversation(
|
||
request: Request,
|
||
session_id: str,
|
||
node_key: str,
|
||
body: ConversationStateUpdateRequest,
|
||
) -> NodeResponse:
|
||
store, user_id, session = await _require_writable_session(request, session_id)
|
||
# 命令式状态机:线程归属校验移入 FOR UPDATE 锁内。
|
||
result = _command_result_or_raise(
|
||
await store.apply_node_command(
|
||
session_id,
|
||
node_key,
|
||
user_id,
|
||
command_id=_command_id_or_new(body.command_id),
|
||
command_type="conversation",
|
||
payload={
|
||
"thread_id": body.thread_id,
|
||
"conversation_snapshot": body.conversation_snapshot,
|
||
},
|
||
)
|
||
)
|
||
return _project_node(result["node"])
|
||
|
||
|
||
@router.post("/sessions/{session_id}/nodes/{node_key}/context")
|
||
async def build_node_context(request: Request, session_id: str, node_key: str) -> dict[str, Any]:
|
||
"""Return the bounded hidden task/intent/upstream context for a node turn."""
|
||
return await _handoff_context(request, session_id, node_key)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/nodes/{node_key}/complete-turn", response_model=NodeResponse)
|
||
async def complete_node_turn(
|
||
request: Request, session_id: str, node_key: str, body: CompleteTurnRequest
|
||
) -> NodeResponse:
|
||
store, user_id, session = await _require_writable_session(request, session_id)
|
||
node = await store.get_node(session_id, node_key, user_id)
|
||
if node is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable node not found")
|
||
# 产物内容持久化是文件 IO,保持在数据库事务之外(避免长事务)。清单为空表示
|
||
# 「本轮无新清单」,由锁内回退到节点当前清单——绝不在事务外用过期快照兜底。
|
||
persisted: list[dict[str, Any]] | None = None
|
||
if body.artifact_manifest:
|
||
persisted = await _persist_artifact_contents(
|
||
node.get("thread_id"),
|
||
body.artifact_manifest,
|
||
)
|
||
# 命令式状态机:能否完成 / 是否新交付 / revision+1 / 失效下游 / 解锁下一 stage,
|
||
# 全部在单一 FOR UPDATE 事务内基于锁定行判定并原子落库(apply_node_command)。
|
||
# 原实现先读快照算好 next_status 再进事务,两个并发请求会基于同一过期快照各自
|
||
# 判定,后写覆盖先写 → 下一 stage 卡 locked / invalidated_by 残留 → 断链。
|
||
result = _command_result_or_raise(
|
||
await store.apply_node_command(
|
||
session_id,
|
||
node_key,
|
||
user_id,
|
||
command_id=_command_id_or_new(body.command_id),
|
||
command_type="complete",
|
||
payload={
|
||
"latest_answer": body.latest_answer,
|
||
"artifact_manifest": persisted,
|
||
"upstream_revision": body.upstream_revision,
|
||
},
|
||
)
|
||
)
|
||
return _project_node(result["node"])
|
||
|
||
|
||
@router.post("/sessions/{session_id}/nodes/{node_key}/reject", response_model=SessionDetailResponse)
|
||
async def reject_node(
|
||
request: Request, session_id: str, node_key: str, body: RejectNodeRequest
|
||
) -> SessionDetailResponse:
|
||
store, user_id, _ = await _require_writable_session(request, session_id)
|
||
# 命令式状态机:可驳回性判定 + 驳回 + 下游失效在同一 FOR UPDATE 事务内完成,
|
||
# 并 bump session/node version。同 command_id 的重复驳回原样重放、不重复计数。
|
||
result = _command_result_or_raise(
|
||
await store.apply_node_command(
|
||
session_id,
|
||
node_key,
|
||
user_id,
|
||
command_id=_command_id_or_new(body.command_id),
|
||
command_type="reject",
|
||
payload={"title": body.title, "reason": body.reason},
|
||
)
|
||
)
|
||
return SessionDetailResponse(
|
||
**result["session"],
|
||
nodes=[_project_node(node) for node in result["nodes"]],
|
||
)
|
||
|
||
|
||
async def _prepare_builtin_run(
|
||
request: Request,
|
||
session_id: str,
|
||
*,
|
||
kind: str,
|
||
follow_up: str = "",
|
||
action_task_id: str | int | None = None,
|
||
action_dictionaries: ActionPlanDictionaries | None = None,
|
||
) -> BuiltinRunPreparationResponse:
|
||
store, user_id, session = await _require_writable_session(request, session_id)
|
||
nodes = await store.list_nodes(session_id, user_id)
|
||
if nodes is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable session not found")
|
||
if not _all_nodes_delivered(nodes, intent_position_id=_session_intent_position_id(session)):
|
||
raise HTTPException(status_code=409, detail="所有岗位产物完成后才能生成或更新内置交付")
|
||
|
||
if kind == "action_plan":
|
||
summary = session.get("summary_snapshot")
|
||
if not isinstance(summary, dict) or summary.get("status") != "done":
|
||
raise HTTPException(status_code=409, detail="请先生成有效的方案总结")
|
||
thread_id = session.get("action_plan_thread_id")
|
||
agent_id = ACTION_PLAN_AGENT_ID
|
||
display_name = "行动规划智能体"
|
||
else:
|
||
if follow_up.strip():
|
||
summary = session.get("summary_snapshot")
|
||
if not isinstance(summary, dict) or summary.get("status") != "done":
|
||
raise HTTPException(status_code=409, detail="请先生成有效的方案总结后再继续问答")
|
||
thread_id = session.get("summary_thread_id")
|
||
agent_id = SUMMARY_AGENT_ID
|
||
display_name = "方案总结智能体"
|
||
|
||
return BuiltinRunPreparationResponse(
|
||
agent_id=agent_id,
|
||
display_name=display_name,
|
||
thread_id=str(thread_id) if isinstance(thread_id, str) and thread_id else None,
|
||
prompt=await _build_builtin_prompt(
|
||
kind=kind,
|
||
session=session,
|
||
nodes=nodes,
|
||
follow_up=follow_up,
|
||
action_task_id=action_task_id,
|
||
action_dictionaries=action_dictionaries,
|
||
),
|
||
source_revisions=_source_revisions(nodes),
|
||
)
|
||
|
||
|
||
async def _complete_builtin_run(
|
||
request: Request,
|
||
session_id: str,
|
||
*,
|
||
kind: str,
|
||
body: CompleteBuiltinRunRequest,
|
||
) -> SessionDetailResponse:
|
||
store, user_id, session = await _require_writable_session(request, session_id)
|
||
nodes = await store.list_nodes(session_id, user_id)
|
||
if nodes is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable session not found")
|
||
if not _all_nodes_delivered(nodes, intent_position_id=_session_intent_position_id(session)):
|
||
raise HTTPException(status_code=409, detail="上游岗位产物已变更或未完成,不能写入旧的内置交付")
|
||
current_revisions = _source_revisions(nodes)
|
||
if body.source_revisions != current_revisions:
|
||
raise HTTPException(status_code=409, detail="岗位产物版本已变化,请基于最新交付重新生成")
|
||
|
||
artifacts = await _persist_artifact_contents(body.thread_id, body.artifact_manifest)
|
||
markdown = [artifact for artifact in artifacts if _is_markdown_artifact(artifact)]
|
||
if not markdown:
|
||
# The report text is already visible in the chat stream. Do not lose
|
||
# that successful result just because the model skipped its final
|
||
# write_file call: mirror the exact response into the thread outputs
|
||
# directory and persist it as the normal report artifact.
|
||
fallback_markdown = await _write_builtin_markdown_fallback(
|
||
thread_id=body.thread_id,
|
||
kind=kind,
|
||
content=body.content,
|
||
)
|
||
fallback_markdown["content"] = body.content
|
||
artifacts.append(fallback_markdown)
|
||
markdown = [fallback_markdown]
|
||
if kind == "action_plan" and not any(_is_json_artifact(artifact) for artifact in artifacts):
|
||
raise HTTPException(status_code=422, detail="行动规划智能体还需交付 action-plan-subtasks.json 子任务清单")
|
||
|
||
snapshot_field = "summary_snapshot" if kind == "summary" else "action_plan_snapshot"
|
||
previous = session.get(snapshot_field)
|
||
previous_snapshot = previous if isinstance(previous, dict) else None
|
||
snapshot = _build_builtin_snapshot(
|
||
kind=kind,
|
||
thread_id=body.thread_id,
|
||
content=body.content,
|
||
artifacts=artifacts,
|
||
source_revisions=current_revisions,
|
||
model=body.model,
|
||
previous=previous_snapshot,
|
||
)
|
||
previous_primary = previous_snapshot.get("primary_artifact") if previous_snapshot else None
|
||
new_primary = snapshot.get("primary_artifact")
|
||
report_changed = not isinstance(previous_primary, dict) or not isinstance(new_primary, dict) or (
|
||
_artifact_signature(previous_primary) != _artifact_signature(new_primary)
|
||
)
|
||
if kind == "summary":
|
||
row = await _update_session_or_raise(
|
||
store,
|
||
session_id,
|
||
user_id,
|
||
summary_thread_id=body.thread_id,
|
||
summary_snapshot=snapshot,
|
||
# A rewritten summary invalidates the derived action plan. A
|
||
# pure Q&A follow-up keeps the existing plan valid.
|
||
action_plan_snapshot=None if report_changed else session.get("action_plan_snapshot"),
|
||
status="active",
|
||
)
|
||
else:
|
||
row = await _update_session_or_raise(
|
||
store,
|
||
session_id,
|
||
user_id,
|
||
action_plan_thread_id=body.thread_id,
|
||
action_plan_snapshot=snapshot,
|
||
status="completed",
|
||
)
|
||
node_rows = await store.list_nodes(session_id, user_id) or []
|
||
return SessionDetailResponse(**row, nodes=[_project_node(node) for node in node_rows])
|
||
|
||
|
||
@router.post("/sessions/{session_id}/summary/prepare", response_model=BuiltinRunPreparationResponse)
|
||
async def prepare_summary(
|
||
request: Request, session_id: str, body: PrepareBuiltinRunRequest
|
||
) -> BuiltinRunPreparationResponse:
|
||
return await _prepare_builtin_run(request, session_id, kind="summary", follow_up=body.message)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/summary/ask", response_model=BuiltinRunPreparationResponse)
|
||
async def ask_summary(
|
||
request: Request, session_id: str, body: PrepareBuiltinRunRequest
|
||
) -> BuiltinRunPreparationResponse:
|
||
return await _prepare_builtin_run(request, session_id, kind="summary", follow_up=body.message)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/summary/complete", response_model=SessionDetailResponse)
|
||
async def complete_summary(
|
||
request: Request, session_id: str, body: CompleteBuiltinRunRequest
|
||
) -> SessionDetailResponse:
|
||
return await _complete_builtin_run(request, session_id, kind="summary", body=body)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/action-plan/prepare", response_model=BuiltinRunPreparationResponse)
|
||
async def prepare_action_plan(
|
||
request: Request, session_id: str, body: PrepareBuiltinRunRequest
|
||
) -> BuiltinRunPreparationResponse:
|
||
return await _prepare_builtin_run(
|
||
request,
|
||
session_id,
|
||
kind="action_plan",
|
||
follow_up=body.message,
|
||
action_task_id=body.task_id,
|
||
action_dictionaries=body.action_dictionaries,
|
||
)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/action-plan/ask", response_model=BuiltinRunPreparationResponse)
|
||
async def ask_action_plan(
|
||
request: Request, session_id: str, body: PrepareBuiltinRunRequest
|
||
) -> BuiltinRunPreparationResponse:
|
||
return await _prepare_builtin_run(
|
||
request,
|
||
session_id,
|
||
kind="action_plan",
|
||
follow_up=body.message,
|
||
action_task_id=body.task_id,
|
||
action_dictionaries=body.action_dictionaries,
|
||
)
|
||
|
||
|
||
@router.post("/sessions/{session_id}/action-plan/complete", response_model=SessionDetailResponse)
|
||
async def complete_action_plan(
|
||
request: Request, session_id: str, body: CompleteBuiltinRunRequest
|
||
) -> SessionDetailResponse:
|
||
return await _complete_builtin_run(request, session_id, kind="action_plan", body=body)
|
||
|
||
|
||
@router.get("/sessions/{session_id}/results", response_model=list[NodeResponse])
|
||
async def list_results(request: Request, session_id: str) -> list[NodeResponse]:
|
||
rows = await _get_store(request).list_results(session_id, _current_user_id(request))
|
||
if rows is None:
|
||
raise HTTPException(status_code=404, detail="Position roundtable session not found")
|
||
return [_project_node(row) for row in rows]
|