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

208 lines
9.4 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.

"""开放接口:按 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,
)