deerflow-code/offline-backend-20260512/backend/app/gateway/routers/position_roundtable.py
2026-09-07 18:24:55 +08:00

1849 lines
78 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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]