208 lines
9.4 KiB
Python
208 lines
9.4 KiB
Python
"""开放接口:按 taskId 直接发起 3qfx(三情分析)的「第一次问答」(无需登录)。
|
||
|
||
外部系统只需把一个 **任务id** 传进来,后端就会:
|
||
|
||
1. 解析 ``3Q`` 业务配置的智能体(``business_mapping`` 里 ``business_code=="3Q"`` 那一行的
|
||
``agent_id``,与前端任务深链 ``goPath=3qfx`` 取的是同一个);
|
||
2. 按 taskId 拉取任务详情(consumer 的 ``cop-task-detail``,地址 / 鉴权头 / mock 都在
|
||
``config.yaml`` 的 ``task_deeplink`` 段),拼出与前端深链一致的开场白
|
||
``任务名称:…\n任务内容:…\n任务id(task_id)为…``;
|
||
3. **新建一条真实会话**(thread 打 ``metadata.{taskId, agent_id}`` 标签,与前端
|
||
``taskThreadMetadata`` / ``useTaskThreads`` 的过滤口径一致),以开场白作为第一条用户消息
|
||
**后台启动一次持久化运行**(写 checkpointer,不阻塞);
|
||
4. **立即返回** ``{thread_id, run_id, agent_id, status}`` —— 调用方拿到「调用成功」即可,
|
||
随后在 3qfx 任务工作区的对话列表里就能看到这条问答并实时观看其进展。
|
||
|
||
该路由前缀 ``/api/open/`` 已在 ``auth_middleware._PUBLIC_PATH_PREFIXES`` 白名单里(且
|
||
``csrf_middleware`` 豁免),**绕过登录鉴权**,以 ``"default"`` 用户桶运行;因 thread 带
|
||
``metadata.taskId``(任务深链对话全局共享、按创建者解析文件路径),任何已登录用户打开该任务
|
||
都能看到同一条会话。
|
||
|
||
请求示例::
|
||
|
||
POST /api/open/3qfx/ask
|
||
{"task_id": "2055"}
|
||
|
||
返回::
|
||
|
||
{"thread_id": "…", "run_id": "…", "agent_id": "…", "status": "running",
|
||
"opening_text": "任务名称:…\n任务内容:…\n任务id(task_id)为2055"}
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import uuid
|
||
from typing import Any
|
||
|
||
import httpx
|
||
from fastapi import APIRouter, HTTPException, Request
|
||
from pydantic import BaseModel, Field
|
||
|
||
from app.gateway.routers.thread_runs import RunCreateRequest
|
||
from app.gateway.services import start_run
|
||
from deerflow.config import get_app_config
|
||
from deerflow.runtime.user_context import DEFAULT_USER_ID, reset_current_user, set_current_user
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
router = APIRouter(prefix="/api/open/3qfx", tags=["open"])
|
||
|
||
# 与前端任务深链 goPath=3qfx → businessCode 的映射保持一致。
|
||
_THREEQ_BUSINESS_CODE = "3Q"
|
||
|
||
|
||
class _OpenDefaultUser:
|
||
"""开放接口(免登录)以「default」用户桶运行的占位 user。
|
||
|
||
满足 :class:`deerflow.runtime.user_context.CurrentUser` 协议(仅需 ``.id``)。
|
||
免登录请求本身没有 user context,而 ``start_run`` 写 thread_meta 时默认 ``user_id=AUTO``,
|
||
无 context 会抛 ``RuntimeError`` → 被 ``start_run`` 的 try/except 当「非致命」吞掉 →
|
||
**thread_meta 行从未创建**,于是前端按 ``metadata.taskId`` 怎么也搜不到这条会话。
|
||
"""
|
||
|
||
id = DEFAULT_USER_ID
|
||
|
||
|
||
class ThreeQAskRequest(BaseModel):
|
||
"""3qfx 开放问答请求体。"""
|
||
|
||
task_id: str = Field(..., description="任务id(cop-task-detail 的 taskId)", min_length=1)
|
||
message: str | None = Field(
|
||
default=None,
|
||
description="可选:自定义首条问题文本。留空(默认)则后端按 task_id 拉取任务详情自动拼开场白。",
|
||
)
|
||
model_name: str | None = Field(default=None, description="可选:覆盖模型(models[] 中的 name)")
|
||
thinking_enabled: bool = Field(default=False, description="是否开启思考;默认关闭")
|
||
|
||
|
||
class ThreeQAskResponse(BaseModel):
|
||
thread_id: str
|
||
run_id: str
|
||
agent_id: str
|
||
status: str
|
||
opening_text: str
|
||
|
||
|
||
class CopTaskDetail(BaseModel):
|
||
"""consumer ``cop-task-detail`` 返回的 ``data`` 字段(与前端 CopTaskDetail 对齐)。"""
|
||
|
||
id: str | int | None = None
|
||
overview: str = ""
|
||
task_direction: str = ""
|
||
task_name: str = ""
|
||
task_content: str = ""
|
||
|
||
|
||
async def _resolve_threeq_agent_id(request: Request) -> str:
|
||
"""取 ``business_code=="3Q"`` 配置的智能体 id;未配置 → 503。"""
|
||
store = getattr(request.app.state, "business_mapping_store", None)
|
||
if store is None:
|
||
raise HTTPException(status_code=503, detail="业务映射未就绪(business_mapping_store 不可用)。")
|
||
mapping = await store.get_mapping(_THREEQ_BUSINESS_CODE)
|
||
agent_id = (mapping or {}).get("agent_id")
|
||
if not agent_id:
|
||
raise HTTPException(status_code=503, detail="3Q(三情分析)业务尚未在「业务链条管理」里配置智能体。")
|
||
return str(agent_id)
|
||
|
||
|
||
def _parse_cop_task_detail(body: Any) -> CopTaskDetail:
|
||
"""从 ``{data: {...}}`` 整包里抽出任务详情(字段缺失按空串)。"""
|
||
data = body.get("data") if isinstance(body, dict) else None
|
||
if not isinstance(data, dict):
|
||
return CopTaskDetail()
|
||
return CopTaskDetail(
|
||
id=data.get("id"),
|
||
overview=data.get("overview") if isinstance(data.get("overview"), str) else "",
|
||
task_direction=data.get("taskDirection") if isinstance(data.get("taskDirection"), str) else "",
|
||
task_name=data.get("taskName") if isinstance(data.get("taskName"), str) else "",
|
||
task_content=data.get("taskContent") if isinstance(data.get("taskContent"), str) else "",
|
||
)
|
||
|
||
|
||
async def _fetch_cop_task_detail(task_id: str) -> CopTaskDetail:
|
||
"""按 taskId 取任务详情。``cop_task_detail_test`` 为真走 mock;否则打真实 consumer 接口。"""
|
||
cfg = get_app_config().task_deeplink
|
||
if cfg.cop_task_detail_test:
|
||
return _parse_cop_task_detail((cfg.mock or {}).get("cop_task_detail"))
|
||
|
||
url = (cfg.cop_task_detail_url or "").strip()
|
||
if not url:
|
||
raise HTTPException(status_code=503, detail="未配置任务详情接口(task_deeplink.cop_task_detail_url),无法按 id 拉取。")
|
||
try:
|
||
async with httpx.AsyncClient(timeout=20.0) as client:
|
||
resp = await client.get(
|
||
url,
|
||
params={"taskId": task_id},
|
||
headers={"Authorization": cfg.cop_task_detail_authorization or "Bearer admin"},
|
||
)
|
||
resp.raise_for_status()
|
||
body = resp.json()
|
||
except Exception as exc: # noqa: BLE001 — 开放接口需把上游错误归一成 502
|
||
logger.warning("拉取 cop-task-detail 失败 (taskId=%s): %s", task_id, exc)
|
||
raise HTTPException(status_code=502, detail=f"拉取任务详情失败:{exc}") from exc
|
||
return _parse_cop_task_detail(body)
|
||
|
||
|
||
def build_opening_text(detail: CopTaskDetail, task_id: str) -> str:
|
||
"""拼 3qfx 开场白,与前端 ``buildOpeningText``(3qfx 分支)逐行对齐。"""
|
||
parts = [
|
||
f"任务名称:{detail.task_name}" if detail.task_name else "",
|
||
f"任务内容:{detail.task_content}" if detail.task_content else "",
|
||
f"任务id(task_id)为{task_id}" if task_id else "",
|
||
]
|
||
return "\n".join(p for p in parts if p)
|
||
|
||
|
||
@router.post("/ask", response_model=ThreeQAskResponse)
|
||
async def ask_3qfx(body: ThreeQAskRequest, request: Request) -> ThreeQAskResponse:
|
||
"""按 taskId 直接发起 3qfx 第一次问答;后台运行,立即返回 thread_id。"""
|
||
task_id = body.task_id.strip()
|
||
if not task_id:
|
||
raise HTTPException(status_code=422, detail="task_id 不能为空。")
|
||
|
||
agent_id = await _resolve_threeq_agent_id(request)
|
||
|
||
# 首条问题:默认按 id 自动拉取任务详情拼开场白;显式传 message 则用它覆盖。
|
||
if body.message and body.message.strip():
|
||
opening_text = body.message.strip()
|
||
else:
|
||
detail = await _fetch_cop_task_detail(task_id)
|
||
opening_text = build_opening_text(detail, task_id)
|
||
if not opening_text:
|
||
raise HTTPException(status_code=502, detail="任务详情为空,无法生成首条问题(可改用 message 显式传入)。")
|
||
|
||
thread_id = str(uuid.uuid4())
|
||
context: dict[str, Any] = {"agent_id": agent_id, "thinking_enabled": body.thinking_enabled}
|
||
if body.model_name:
|
||
context["model_name"] = body.model_name
|
||
|
||
run_body = RunCreateRequest(
|
||
assistant_id=agent_id,
|
||
input={"messages": [{"role": "user", "content": opening_text}]},
|
||
# taskId + agent_id 让该会话被 3qfx 任务工作区的 useTaskThreads 过滤命中(按任务+业务隔离),
|
||
# 且 metadata.taskId 触发后端「任务深链对话全局共享 / 按创建者解析文件」语义。
|
||
metadata={"taskId": task_id, "agent_id": agent_id},
|
||
context=context,
|
||
multitask_strategy="reject",
|
||
)
|
||
|
||
# 免登录 → 无 user context。显式设成「default」用户桶,否则 start_run 写 thread_meta 时
|
||
# resolve_user_id(AUTO) 抛 RuntimeError 被吞掉,thread_meta 建不出来、前端按 taskId 搜不到。
|
||
# run_agent 是后台 task,在 create_task 时快照本 context,故后台运行也归属 default 用户
|
||
# (与 metadata.taskId「任务深链对话全局共享、按创建者解析文件」语义一致)。
|
||
# enforce_access=False:开放接口已自行解析 admin 配置的 3Q 智能体,不再按当前(无)登录用户做可见性校验。
|
||
token = set_current_user(_OpenDefaultUser())
|
||
try:
|
||
record = await start_run(run_body, thread_id, request, enforce_access=False)
|
||
finally:
|
||
reset_current_user(token)
|
||
|
||
return ThreeQAskResponse(
|
||
thread_id=record.thread_id,
|
||
run_id=record.run_id,
|
||
agent_id=agent_id,
|
||
status=record.status.value,
|
||
opening_text=opening_text,
|
||
)
|