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

908 lines
44 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.

"""圆桌规划 · 第一步「任务理解」Gateway 路由。
前端页面:``/page/workspace/roundtable/planning`` Step 1
前端模块:``frontend-web/src/roundtable-planning/api/intent.ts``
职责
----
在独立 thread 上运行内置智能体 ``roundtable-intent``(SOUL 见
``.deer-flow/agents/roundtable-intent/SOUL.md``),与用户进行 2–5 轮澄清,
最终输出结构化意图 ``{objective, constraints, assumptions}``,供 Step 2
总控协调智能体作为初始 prompt。
与 ``multi_agent.py`` 的关系
---------------------------
- **复用**同一套 loopback 工具:``_stream_upstream``、``_run_payload``、
``_parse_last_messages`` 等,保证 SSE 透传 + 末尾状态帧的解析逻辑一致。
- **不复用** ``/init`` 的协调智能体创建逻辑:intent 是长期注册的内置 agent,
只需 ``POST /api/intent/init`` 创建 thread 即可。
SSE 协议(与前端 ``streamIntent`` 对齐)
--------------------------------------
1. 透传 LangGraph ``/api/threads/{id}/runs/stream`` 的全部 ``data:`` 行
(含 ``messages-tuple`` 增量,供 UI 流式打字)。
2. 流结束后追加**一条** JSON 状态帧(``data: {...}\\n\\n``),字段:
``status`` — 终态枚举:
- ``asking``:本轮无 ``ask_clarification``,模型仍在追问或闲聊;
- ``clarification``:本轮以澄清工具结束,需用户选选项/自定义输入;
- ``done``:已解析 ``[INTENT_READY]`` + JSON 块,``summary`` 有值;
- ``error``:上游或网关异常。
``content`` — 展示用 AI 正文(澄清时只保留 tool 调用前的简短引导;
tool 后的尾随文本会被丢弃,选项发出后立即等待用户选择)。
仅 ``status == "done"`` 时带 ``summary``;澄清时还带
``question``、``clarification_type``、``clarification_context``、
``options``、``allow_custom``、``allow_multiple``(见 ``resolve_allow_multiple``)。
多卡回合还带 ``clarifications`` 完整数组;顶层字段始终镜像第一张卡以兼容旧客户端。
终态判定优先级(实现细节,勿随意调整)
------------------------------------
1. **``[INTENT_READY]`` 优先**:从最新消息向前扫描 AIMessage,正则匹配
SOUL 规定的 ``[INTENT_READY]\\n```json {...} `````;即使最后一条消息仍带
过期的 ``ask_clarification`` tool_call,也返回 ``done``(避免前端卡在澄清 UI)。
2. 否则看最后一条带 ``tool_calls`` 的 AIMessage 是否含 ``ask_clarification``。
3. 否则 ``asking``。
工具白名单
----------
``excluded_tools`` 屏蔽除 ``ask_clarification`` 外的一切工具,防止模型偏离
澄清任务去写文件、派活等(SOUL 文字约束不可靠,硬排除更稳)。
"""
from __future__ import annotations
import hashlib
import json
import logging
import os
import re
import uuid
from collections.abc import AsyncIterator
from datetime import UTC, datetime, timedelta
from typing import Any
import httpx
from fastapi import APIRouter, Body, HTTPException, Request
from fastapi.responses import StreamingResponse
from app.gateway.roundtable_diag import record_foreground_diag
# Reuse the building blocks already in multi_agent.py — keeps both routers
# identical in connection handling, payload shape, and error mapping.
from app.gateway.routers._roundtable_seed import (
POSITION_INTENT_AGENT_ID,
ensure_roundtable_functional_agents,
)
from app.gateway.routers.multi_agent import (
_ROUNDTABLE_THREAD_METADATA,
_auth_headers,
_create_thread_direct,
_diagnose_exception,
_find_dispatch_message,
_find_trailing_ai_text,
_flatten_content,
_loopback_base,
_make_loopback_client,
_parse_last_messages,
_run_payload,
_stream_upstream,
)
from deerflow.tools.builtins.clarification_utils import resolve_allow_multiple
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/intent", tags=["intent"])
# ── Phase 4:意图回合跨 worker 单飞 + 幂等 ─────────────────────────────────────
# 进程内 RunManager 是每 worker 内存态,多 worker 下同 thread 的两个请求落到不同
# worker 会各启一个 run(无 409)。单飞只能靠 DB:intent_turns 表对 (thread_id,
# client_turn_id) 的唯一键 + 条件 UPDATE 抢占。前端每回合生成稳定 client_turn_id,
# 重试 / 双击复用,杜绝「同一回合被启两个 run → 双发意图 / 重复澄清卡」。
_WORKER_ID = f"w-{os.getpid()}-{uuid.uuid4().hex[:8]}"
_INTENT_TURN_LEASE_TTL_SECONDS = 120
def _intent_message_hash(message: str, files: list, task_context: dict | None) -> str:
"""回合内容指纹:message + files + task_context 的 sha256。
检测「同一 client_turn_id 但内容变了」的客户端 bug(→ conflict)。文件按
(filename, path) 归一,避免临时字段(size/status)造成假性不一致。
"""
file_repr = sorted(
(str(f.get("filename") or ""), str(f.get("path") or "")) for f in (files or [])
)
blob = json.dumps(
{"message": message, "files": file_repr, "task_context": task_context or {}},
ensure_ascii=False,
sort_keys=True,
)
return hashlib.sha256(blob.encode("utf-8")).hexdigest()
def _get_intent_turn_store(request: Request):
"""意图回合单飞 store(未配置 DB 时为 None → 跳过单飞,走原流程)。"""
return getattr(request.app.state, "intent_turn_store", None)
# 内置 agent 目录名,须与 ``.deer-flow/agents/roundtable-intent/`` 一致。
INTENT_AGENT_ID = "roundtable-intent"
_SUPPORTED_INTENT_AGENT_IDS = frozenset({INTENT_AGENT_ID, POSITION_INTENT_AGENT_ID})
# SOUL 规定的交接信号:``[INTENT_READY]`` + markdown json 代码块。
# group(1) 解析为 IntentSummary;解析失败则继续走澄清/asking 分支。
_INTENT_READY_PATTERN = re.compile(
r"\[INTENT_READY\]\s*```json\s*(\{.*?\})\s*```",
re.DOTALL,
)
def _requested_intent_agent_id(payload: dict[str, Any]) -> str:
"""Resolve the explicitly supported Step-1 agent without opening a route
for arbitrary custom-agent execution.
Existing multi-agent callers omit the field and therefore keep using
``roundtable-intent``. Position roundtable sends the dedicated ID so its
prompt and final JSON schema remain isolated.
"""
requested = str(payload.get("intent_agent_id") or "").strip()
if not requested:
return INTENT_AGENT_ID
if requested not in _SUPPORTED_INTENT_AGENT_IDS:
raise HTTPException(status_code=400, detail="Unsupported intent agent")
return requested
def _clarification_blocks_ready(
intent_agent_id: str,
clarification_calls: Any,
*,
is_fresh: bool,
) -> bool:
"""A fresh clarification always pauses intent resolution.
A tool call and a premature ``[INTENT_READY]`` block may coexist in the
same upstream run. The tool call is an explicit request for user input,
so it wins for every Step-1 workspace (including 8BF): do not consume the
assistance card before the user can answer it. Stale calls from earlier
turns remain harmless because only a fresh call blocks READY.
``intent_agent_id`` stays in the signature for route-call compatibility
and to keep the policy extensible per agent if a future flow needs one.
"""
del intent_agent_id
return bool(clarification_calls) and is_fresh
def _intent_thread_metadata(
task_context: dict[str, Any] | None,
*,
intent_agent_id: str = INTENT_AGENT_ID,
) -> dict[str, Any]:
"""Build task-scoped metadata for a newly-created intent thread.
Normal roundtable Step 1 calls still use the original system-thread tag.
Position roundtable additionally supplies the task card snapshot, so the
intent dialogue can be indexed and restored before the user confirms a
business chain.
"""
metadata = dict(_ROUNDTABLE_THREAD_METADATA)
if not isinstance(task_context, dict):
return metadata
task = task_context.get("task")
task_id = str(task.get("id") or "").strip() if isinstance(task, dict) else ""
if not task_id:
return metadata
metadata["taskId"] = task_id
# 8BF 使用通用意图接口但有明确的上下文 kind;它不能被岗位会商的
# 会话索引误认为岗位任务,否则两个页面的历史记录会串在一起。
# 其它历史调用保持原有 metadata 形状,避免影响既有会商恢复逻辑。
is_eightbf_context = task_context.get("kind") == "eightbf_roundtable_task"
if intent_agent_id == POSITION_INTENT_AGENT_ID or not is_eightbf_context:
metadata.update(
{
"position_roundtable_task_id": task_id,
"position_task_id": f"{task_id}:intelligence",
"position_id": "intelligence",
"position_roundtable_phase": "intent",
}
)
if intent_agent_id == POSITION_INTENT_AGENT_ID:
metadata["position_roundtable_intent_agent_id"] = intent_agent_id
return metadata
def _normalize_question(text: str) -> str:
"""归一化澄清问题文本,用于"和过往是否重复"的判断。
LLM 复述时常见模式:
- 完全一样的字符串
- 仅前后空白、标点、emoji 不同(如多/少一个问号)
- 仅大小写差异(中文场景几乎不见,但兼容英文/拼音)
所以归一化时统一:小写 + 去除非字母/数字/中文字符 + 折叠空白。
Levenshtein/语义相似度等模糊匹配先不做,等观察到误判再迭代。
"""
if not isinstance(text, str):
return ""
lowered = text.strip().lower()
# 只保留中英文字符与数字,去掉 ?。!?,。"":/、() 等标点 + emoji + 空白。
return re.sub(r"[^0-9a-z一-鿿]+", "", lowered)
def _normalize_options(raw: object) -> tuple[str, ...]:
"""把 options 拍平成可比较的归一化元组,用于"选项是否和上一轮一致"的判断。"""
if raw is None:
return ()
if isinstance(raw, str):
# 偶有模型把 list 序列化成 JSON 字符串发回来。
try:
raw = json.loads(raw)
except (json.JSONDecodeError, TypeError):
return (_normalize_question(raw),)
if not isinstance(raw, list):
return (_normalize_question(str(raw)),)
out: list[str] = []
for item in raw:
if isinstance(item, dict):
label = item.get("label") or item.get("title") or item.get("name") or item.get("value") or item.get("id") or ""
out.append(_normalize_question(str(label)))
else:
out.append(_normalize_question(str(item)))
# 顺序无关:LLM 偶尔会调换选项顺序,但语义上还是同一组。
return tuple(sorted(o for o in out if o))
def _count_past_clarifications(messages: list[dict[str, Any]]) -> int:
"""统计当前 thread 内已经发出过的 ``ask_clarification`` 次数。
用来给 ``stream_intent`` 做"3 轮硬上限"兜底:
- 0-1 次:正常流程,模型自行决定继续问还是收敛
- 2 次(即将进入第 3 轮):prompt 追加系统提示,要求本轮必须 INTENT_READY
- ≥3 次:工具白名单移除 ``ask_clarification``,模型物理上无法再问
只数"已经成型"的 tool_call,不区分用户是否回答过——LangGraph 把 ask
工具的调用永久记录在 AIMessage.tool_calls 里,数它即可。
"""
if not isinstance(messages, list):
return 0
count = 0
for msg in messages:
if not isinstance(msg, dict):
continue
if (msg.get("type") or "").lower() not in ("ai", "aimessage", "aimessagechunk"):
continue
for tc in msg.get("tool_calls") or []:
if (tc.get("name") or "").lower() == "ask_clarification":
count += 1
return count
def _count_past_clarification_rounds(messages: list[dict[str, Any]]) -> int:
"""Count completed assistance *rounds*, rather than question cards.
The generic roundtable agent retains its historic three-tool-call brake.
Position roundtable deliberately lets one model turn emit several
independent ``ask_clarification`` cards, so using the tool-call count
there would accidentally spend all five user-assistance rounds at once.
One AI message containing one or more clarification calls is one round.
"""
if not isinstance(messages, list):
return 0
rounds = 0
for msg in messages:
if not isinstance(msg, dict):
continue
if (msg.get("type") or "").lower() not in ("ai", "aimessage", "aimessagechunk"):
continue
if any(
(tc.get("name") or "").lower() == "ask_clarification"
for tc in msg.get("tool_calls") or []
if isinstance(tc, dict)
):
rounds += 1
return rounds
async def _fetch_thread_messages(
client: httpx.AsyncClient,
base: str,
headers: dict[str, str],
thread_id: str,
) -> list[dict[str, Any]]:
"""读取当前 thread checkpoint 的 messages 列表(best-effort)。
用 ASGITransport loopback 调 ``GET /api/threads/{tid}/state``,失败时返回
空列表——计数兜底是个"刹车",失败时退化成原有的纯 SOUL 自律行为,不阻断主流程。
"""
try:
resp = await client.get(f"{base}/api/threads/{thread_id}/state", headers=headers)
except Exception as exc:
logger.warning("intent: get thread state failed (%s) — falling back to soul-only flow", exc)
return []
if resp.status_code != 200:
logger.warning(
"intent: get thread state %s returned %s — falling back to soul-only flow",
thread_id,
resp.status_code,
)
return []
try:
payload = resp.json()
except Exception:
return []
values = payload.get("values") if isinstance(payload, dict) else None
if not isinstance(values, dict):
return []
messages = values.get("messages")
return messages if isinstance(messages, list) else []
# 硬上限:超过这个轮数,无论模型怎么想,都强制移除 ask_clarification 工具。
# 通用会商沿用 3 次工具调用的历史保护;岗位会商按"一轮可多张卡"计数,
# 允许最多 5 次用户协助。
_MAX_CLARIFICATION_ROUNDS = 3
_POSITION_MAX_CLARIFICATION_ROUNDS = 5
def _max_clarification_rounds(intent_agent_id: str) -> int:
return (
_POSITION_MAX_CLARIFICATION_ROUNDS
if intent_agent_id == POSITION_INTENT_AGENT_ID
else _MAX_CLARIFICATION_ROUNDS
)
def _is_repeated_clarification(
messages: list[dict[str, Any]],
*,
new_question: str,
new_options: object,
new_tool_call_id: str | None = None,
) -> bool:
"""检测本轮的 ask_clarification 是否在过往轮次问过。
判定为重复必须同时满足:
1. 历史里至少有一次 ``ask_clarification`` 的问题文本 / 选项 与本次归一化相同;
2. 用户在那次之后**已经发过消息**(即"已回答过")。
没有这条限制会把"模型在同一 run 内多次输出同一 tool_call_chunks"
误判成重复,导致首次澄清都吞掉。
3. 不是同一个 tool_call(LangGraph 偶有把同一 call 拆成多块返回的情况)。
传回 True → 上层应把 status 从 ``clarification`` 降级为 ``asking``,
避免前端弹同一张卡片(配合 SOUL §2.5 prompt 双层防御)。
"""
norm_new_q = _normalize_question(new_question)
norm_new_opts = _normalize_options(new_options)
if not norm_new_q and not norm_new_opts:
return False
seen_match_index: int | None = None
for idx, msg in enumerate(messages):
if not isinstance(msg, dict):
continue
if (msg.get("type") or "").lower() not in ("ai", "aimessage", "aimessagechunk"):
continue
for tc in msg.get("tool_calls") or []:
if (tc.get("name") or "").lower() != "ask_clarification":
continue
if new_tool_call_id and tc.get("id") == new_tool_call_id:
# 同一个 tool_call 的不同 chunk,不算重复。
continue
args = tc.get("args") or {}
past_q = _normalize_question(str(args.get("question") or ""))
past_opts = _normalize_options(args.get("options"))
# 问题文本和选项任一明确重叠就视为重复:
# - question 相同(options 可能空)
# - options 相同且非空(question 可能被改写)
q_match = bool(past_q) and past_q == norm_new_q
opt_match = bool(past_opts) and past_opts == norm_new_opts
if q_match or opt_match:
seen_match_index = idx
break
if seen_match_index is not None:
break
if seen_match_index is None:
return False
# 校验用户在该次澄清之后是否已经发过 human 消息(= 回答过)。
for msg in messages[seen_match_index + 1:]:
if not isinstance(msg, dict):
continue
msg_type = (msg.get("type") or "").lower()
if msg_type in ("human", "humanmessage", "user"):
return True
return False
@router.post("/init")
async def init_intent(request: Request, payload: dict = Body(...)) -> dict[str, Any]:
"""创建任务理解会话用的空 thread。
请求体:``{ "model": "<可选>" }``(可选,当前 init 未把 model 写入 thread,
实际模型在 ``/stream`` 的每次 run 里指定;缺省由 lead_agent 回退到
``config.yaml`` 的 ``models[0]``,前端通常传当前下拉框选中的 model 名)。
响应::
{ "status": "success", "agent_id": "roundtable-intent", "thread_id": "<uuid>" }
与 ``POST /api/multi-agent/init`` 的区别:不创建协调智能体、不建多席位
thread 映射;agent 本身已通过磁盘 ``.deer-flow/agents/`` 同步注册。
"""
# 调用前自检:内网部署 .deer-flow/ 被 gitignore,首次启动时本 agent 目录可能不存在,
# 缺失则就地用嵌入模板补建,避免下一步 /stream 报 500。
ensure_roundtable_functional_agents()
intent_agent_id = _requested_intent_agent_id(payload)
# Keep Step-1 initialization aligned with multi-agent initialization:
# create the empty thread in-process rather than looping back through
# 127.0.0.1:<gateway-port>. Intranet proxies and port mappings can make
# that loopback fail only for automatic page initialization. Streaming
# below still uses the TCP loopback so SSE remains incrementally flushed.
raw_task_context = payload.get("task_context")
task_context = raw_task_context if isinstance(raw_task_context, dict) else None
try:
thread_id = await _create_thread_direct(
request,
metadata=_intent_thread_metadata(task_context, intent_agent_id=intent_agent_id),
)
except HTTPException:
# 上游已用结构化方式失败,原样传给前端不需要二次包装。
raise
except Exception as exc:
# 把根因暴露给前端而不是吞成裸 500——内网部署常见的「config.yaml 模型外网不通」
# / 「内置 agent 目录缺失」/ 「loopback 走代理」都能在前端面板直接看到。
# FastAPI 把 detail=dict 序列化为 {"detail": {...}};前端 initIntent
# 在 !res.ok 分支里读 detail.hint / trace_tail 渲染诊断面板。
diagnostic = _diagnose_exception(exc, context="intent_init")
logger.error("intent init failed: %s", diagnostic, exc_info=exc)
await record_foreground_diag(
request, stage="step1_intent", level="error", event="intent_init_failed",
message=f"第一步意图初始化失败:{diagnostic.get('message') or exc}",
detail=diagnostic,
)
raise HTTPException(status_code=500, detail=diagnostic) from exc
return {
"status": "success",
"agent_id": intent_agent_id,
"thread_id": thread_id,
}
@router.post("/stream")
async def stream_intent(request: Request, payload: dict = Body(...)) -> StreamingResponse:
"""Drive one turn of the intent dialogue via SSE.
Body: ``{thread_id, message, model?}``. Streams upstream LangGraph frames
through transparently, then appends one status frame:
{
"status": "asking" | "clarification" | "done" | "error",
"content": "<final AI text>",
"summary": {...} | null, // only when status == "done"
// when status == "clarification":
"question", "clarification_type", "clarification_context",
"options", "allow_custom", "allow_multiple"
}
"""
thread_id = str(payload.get("thread_id") or "").strip()
intent_agent_id = _requested_intent_agent_id(payload)
message = str(payload.get("message") or "").strip()
# 不写死任何默认模型名:空串会让 lead_agent._resolve_model_name() 回退到
# config.yaml 的 models[0],与 /api/ai-writing 等其它路由一致,内网部署只配
# 自己的本地模型即可工作。
model_name = str(payload.get("model") or "")
# 用户在 Step 1 上传的参考文件(已经过 /api/threads/{thread_id}/uploads)。
# 前端传 [{filename, size, path, status}, ...],这里原样透传给 _run_payload →
# additional_kwargs.files → UploadsMiddleware 自动注入清单 + 解锁 read_file 等工具。
raw_files = payload.get("files")
files = raw_files if isinstance(raw_files, list) else []
raw_task_context = payload.get("task_context")
task_context = raw_task_context if isinstance(raw_task_context, dict) else None
if not thread_id:
raise HTTPException(status_code=400, detail="thread_id is required")
# 允许「仅文件、无文本」的回合:带了文件就不强制要求 message。
if not message and not files:
raise HTTPException(status_code=400, detail="message is required")
# 调用前自检:防御性二次检查——若管理员在 init 之后手动删了目录,这里也能自愈。
ensure_roundtable_functional_agents()
base = _loopback_base(request)
headers = _auth_headers(request)
fallback_instruction = (
"四个岗位会商字段都必须保留为字符串数组,依据已获得的用户确认形成明确结论,"
"不得写待确认、未明确、需补充、暂无或缺少"
if intent_agent_id == POSITION_INTENT_AGENT_ID
else "缺的字段用 assumptions 兜底"
)
async def stream_generator() -> AsyncIterator[str]:
# ASGITransport:进程内直调本机 ASGI app,**不走 TCP / 不经过网络栈**,
# 因此既无须 trust_env=False 防代理,也不存在 loopback 40s 超时这种坑。
async with _make_loopback_client(request) as client:
try:
# —— Phase 4:意图回合跨 worker 单飞 + 幂等 ——
# 同一 (thread_id, client_turn_id) 至多一个 run:重复 / 双击 / 网络重发
# 命中既有行即幂等返回,绝不重复启 run(杜绝双发意图 / 重复澄清卡)。
turn_store = _get_intent_turn_store(request)
client_turn_id = str(payload.get("client_turn_id") or "").strip() or uuid.uuid4().hex
msg_hash = _intent_message_hash(message, files, task_context)
claim_action = "run" # 无 store 时走原流程
if turn_store is not None:
claim = await turn_store.try_claim(
thread_id=thread_id,
client_turn_id=client_turn_id,
message_hash=msg_hash,
lease_owner=_WORKER_ID,
lease_until=datetime.now(UTC)
+ timedelta(seconds=_INTENT_TURN_LEASE_TTL_SECONDS),
)
claim_action = claim.get("action") or "run"
if claim_action == "replay":
cached = (claim.get("row") or {}).get("result_json")
if isinstance(cached, dict):
yield f"data: {json.dumps(cached, ensure_ascii=False)}\n\n"
return
# 缓存缺失/损坏 → 降级重跑(极少见)。
claim_action = "run"
elif claim_action == "duplicate":
# 既有回合仍在跑且租约有效:本请求是重复 in-flight,丢弃
# (原请求正在流式吐结果)。前端按 duplicate 标记忽略。
dup = {"status": "asking", "content": "", "summary": None, "duplicate": True}
yield f"data: {json.dumps(dup, ensure_ascii=False)}\n\n"
return
elif claim_action == "conflict":
err_conflict = {
"status": "error",
"content": "意图回合幂等键冲突:同一 client_turn_id 但消息内容不同。",
"summary": None,
}
yield f"data: {json.dumps(err_conflict, ensure_ascii=False)}\n\n"
return
async def _record_terminal(frame: dict[str, Any]) -> None:
"""终态帧落库(仅本 worker 赢得领取时),供后续重复请求幂等回放。"""
if turn_store is None or claim_action != "run":
return
try:
await turn_store.record_result(
thread_id=thread_id,
client_turn_id=client_turn_id,
lease_owner=_WORKER_ID,
status="error" if frame.get("status") == "error" else "done",
result_json=frame,
)
except Exception:
logger.warning("intent_turn record_result failed", exc_info=True)
# —— 澄清轮次硬上限兜底 ——
# 通用会商保持原来的 3 次 ask 工具调用上限;岗位会商允许一个
# AI turn 同时产生多张独立协助卡,所以必须按 AI 消息(一次
# "提问 + 用户回答")计轮,不能按卡片数计轮。
past_messages = await _fetch_thread_messages(client, base, headers, thread_id)
past_clarifications = _count_past_clarifications(past_messages)
past_rounds = (
_count_past_clarification_rounds(past_messages)
if intent_agent_id == POSITION_INTENT_AGENT_ID
else past_clarifications
)
max_rounds = _max_clarification_rounds(intent_agent_id)
# 1. ask_clarification 的工具白名单——超过硬上限就移除。
base_excluded_tools = [
"web_search",
"present_files",
"view_image",
"agent_orchestration",
"bash",
"ls",
"read_file",
"write_file",
"str_replace",
"memory",
"hindsight_recall",
"hindsight_reflect",
"hindsight_retain",
"task",
"write_todos",
]
if past_rounds >= max_rounds:
base_excluded_tools.append("ask_clarification")
logger.info(
"intent stream: past_rounds=%d >= %d, force-disabling "
"ask_clarification to break the loop",
past_rounds,
max_rounds,
)
# 带了上传文件时,解锁「读」类工具(read_file / ls / grep / view_image),
# 否则 UploadsMiddleware 注入了文件清单、模型却没工具读取。仍保留对
# write_file / str_replace / present_files / bash / agent_orchestration 的封锁,
# 防止意图阶段越界做方案/派活。
if files:
read_tools = {"read_file", "ls", "grep", "view_image"}
base_excluded_tools = [t for t in base_excluded_tools if t not in read_tools]
# 2. message 前缀——给模型一个明确的"本轮必须收敛"信号。
# 仅文件、无文本时给一句兜底引导,让模型先读文件再澄清意图。
effective_message = message or "我上传了文件,请先阅读文件内容,再据此理解并澄清我的任务意图。"
if past_rounds < max_rounds:
effective_message = (
"[澄清输出硬约束] 如果本轮调用 ask_clarification,"
"工具调用完成后必须立即停止输出并等待用户选择;"
"禁止补充说明、复述选项、催促选择、追加问题或输出 "
"[INTENT_READY]。需要确认的全部内容只能写入工具参数。\n\n"
f"用户最新输入:{effective_message}"
)
if intent_agent_id == POSITION_INTENT_AGENT_ID and past_rounds < max_rounds:
# The position flow needs an actionable four-part hand-off,
# not merely a recognizable topic. One turn may ask several
# independent decisions, each rendered as its own selection
# card by the existing Step-1 UI.
effective_message = (
"[岗位会商收敛检查] 仅识别出主题不等于意图已明确。"
"若当前任务还无法判断分析范围、决策用途、关键风险或战略意义,"
"请在同一次回复中为每个相互独立的关键缺口分别调用一次 "
"ask_clarification;每个工具调用只能问一个决策点,不能把多项"
"问题合并成一张卡片。用户会一次性回答这些卡片。"
"如果四项意图已经足够明确,则直接输出最终结果。"
"一旦调用完本轮所有 ask_clarification,必须立即停止输出,"
"不要再补充说明、复述选项、催促选择或继续提问,等待用户选择。"
"不要以风险提示和战略意义为空的 [INTENT_READY] 结束。"
"本轮是第 %d/%d 轮用户协助;第 %d 轮若仍有缺口,必须把所有"
"剩余关键缺口分别发成卡片,收到回答后立即收敛。\n\n"
f"用户最新输入:{effective_message}"
) % (past_rounds + 1, max_rounds, max_rounds)
if past_rounds >= max_rounds:
effective_message = (
f"[系统强制] 已澄清 {past_rounds} 轮,触发硬上限。"
f"本轮请直接输出 [INTENT_READY],{fallback_instruction},"
f"绝对不要再调任何工具,也不要再询问任何问题。\n\n"
f"用户最新输入:{message}"
)
elif (
intent_agent_id != POSITION_INTENT_AGENT_ID
and past_rounds >= max_rounds - 1
):
effective_message = (
f"[系统提示] 已澄清 {past_rounds} 轮,这是最后一轮。"
f"本轮请直接输出 [INTENT_READY],{fallback_instruction}。\n\n"
f"用户最新输入:{message}"
)
body = _run_payload(
agent_name=intent_agent_id,
model_name=model_name,
thread_id=thread_id,
new_message=effective_message,
# Intent agent only ever needs ask_clarification (and even that
# is yanked once past_clarifications hits the cap above). Block
# everything else so the model can't drift into solutioning
# (e.g. saving a design doc with write_file) — SOUL prose is
# not reliable on its own, hard exclusion is.
excluded_tools=base_excluded_tools,
skill_stop_names=[],
extra_context={"task_external_context": task_context} if task_context else None,
files=files,
)
# captured:上游 SSE 全部 data 行,用于流结束后从 values 快照取 messages。
captured: list[str] = []
async for line, done in _stream_upstream(client, base, headers, thread_id, body):
if line is not None:
yield line if line.endswith("\n") else line + "\n"
elif done is not None:
captured = done
messages = _parse_last_messages(captured)
current_rounds = (
_count_past_clarification_rounds(messages)
if intent_agent_id == POSITION_INTENT_AGENT_ID
else _count_past_clarifications(messages)
)
clarification_is_fresh = current_rounds > past_rounds
# A fresh clarification is an explicit pause. Some models emit
# a premature [INTENT_READY] block after its tool call in the
# same upstream run; the user-assistance card must win so the
# user can answer it before any summary is accepted.
dispatch_msg = _find_dispatch_message(messages)
tool_calls = dispatch_msg.get("tool_calls") or []
clarification_calls = [
tc
for tc in tool_calls
if isinstance(tc, dict)
and (tc.get("name") or "").lower() == "ask_clarification"
]
# --- 终态分支 1:结构化意图已就绪(无新追问时)---
# 模型可能在较早的 AIMessage 里已输出 [INTENT_READY],而最后一轮
# 仍残留 ask_clarification 的 tool_calls;若它是本轮新产生的追问,
# 必须先等待用户回答,不能让 READY 覆盖可交互的协助卡。
summary: dict[str, Any] | None = None
ready_text: str = ""
for msg in reversed(messages):
if not isinstance(msg, dict):
continue
if (msg.get("type") or "").lower() not in ("ai", "aimessage", "aimessagechunk"):
continue
msg_text = _flatten_content(msg.get("content", ""))
m = _INTENT_READY_PATTERN.search(msg_text)
if m:
try:
summary = json.loads(m.group(1))
except json.JSONDecodeError as exc:
logger.warning("intent summary JSON parse failed: %s", exc)
summary = None
# SOUL §0 硬规则 3 要求 ```json``` 反引号闭合后不得再输出任何字符,
# 但模型仍偶尔会在后面接 "那么我帮你列一下大纲……" 之类的尾巴。
# 终态帧只回到 marker 块结尾 + 一个换行,把废话从 content 里裁掉——
# 草稿持久化 / 跨设备重放 / 前端 splitIntentReadyBlock 的 after
# 段都因此干净。流式期间用户仍会逐 token 看到那段废话,前端
# Step1Panel 已经把 `parts.after` 隐藏掉,体验闭环。
ready_text = msg_text[: m.end()].rstrip() + "\n"
break
if (
summary
and ready_text
and not _clarification_blocks_ready(
intent_agent_id,
clarification_calls,
is_fresh=clarification_is_fresh,
)
):
final_frame = {
"status": "done",
"content": ready_text,
"summary": summary,
}
await _record_terminal(final_frame)
yield f"data: {json.dumps(final_frame, ensure_ascii=False)}\n\n"
return
# --- 终态分支 2:本轮以 ask_clarification 结束 ---
# _find_dispatch_message:从后往前找带 tool_calls 的 AIMessage
# (SkillStop/澄清 ToolMessage 之后的那条 AI 消息)。
content = _flatten_content(dispatch_msg.get("content", ""))
if clarification_calls:
clarification_payloads: list[dict[str, Any]] = []
for clarification_call in clarification_calls:
cargs = clarification_call.get("args") or {}
question = str(cargs.get("question") or "").strip()
# 逐卡过滤已被用户回答过的重复问题,不能因为其中一张
# 重复卡而吞掉同轮其它有效的协助卡。
if _is_repeated_clarification(
messages,
new_question=question,
new_options=cargs.get("options"),
new_tool_call_id=clarification_call.get("id"),
):
logger.warning(
"intent stream: dropping duplicate ask_clarification "
"(question=%r) — already asked & answered earlier in thread",
question[:80],
)
continue
clarification_payloads.append(
{
"id": str(clarification_call.get("id") or ""),
"question": question,
"clarification_type": str(cargs.get("clarification_type") or ""),
"clarification_context": str(cargs.get("context") or ""),
"options": cargs.get("options") or [],
"allow_custom": bool(cargs.get("allow_custom", True)),
"allow_multiple": resolve_allow_multiple(cargs),
}
)
if not clarification_payloads:
# 这里只在所有卡片都因重复而被过滤时才允许使用 tool 后文本;
# 此时前端没有选项需要等待,不属于澄清暂停态。
follow_up = _find_trailing_ai_text(messages, exclude_text=content)
fallback_msg = (
follow_up
or content
or "(系统识别到智能体重复了上一问题,已自动跳过。请补充其他信息继续推进。)"
)
final_frame = {
"status": "asking",
"content": fallback_msg,
"summary": None,
}
await _record_terminal(final_frame)
yield f"data: {json.dumps(final_frame, ensure_ascii=False)}\n\n"
return
first = clarification_payloads[0]
# ask_clarification 是严格的回合终点。部分模型会在工具调用后
# 再生成“请从以上选项中选择……”等尾随 AIMessage;不能把它
# 返回给前端,否则用户看到卡片后还会继续冒出一段话。
# 仅保留工具调用自身携带的前置短句;没有则直接由卡片展示问题。
display_content = content
final_frame: dict[str, Any] = {
"status": "clarification",
"content": display_content,
# 兼容旧客户端保留首张卡的顶层字段;完整列表使得网络
# 重连、缓存重放等没有实时 tool frame 的路径也能还原多卡。
**first,
"clarifications": clarification_payloads,
"summary": None,
}
await _record_terminal(final_frame)
yield f"data: {json.dumps(final_frame, ensure_ascii=False)}\n\n"
return
# --- 终态分支 3:普通 AI 回复,继续多轮对话 ---
final_frame: dict[str, Any] = {
"status": "asking",
"content": content,
"summary": None,
}
await _record_terminal(final_frame)
yield f"data: {json.dumps(final_frame, ensure_ascii=False)}\n\n"
except HTTPException as exc:
logger.warning("intent stream HTTPException: %s %s", exc.status_code, exc.detail)
# HTTPException 的 detail 可能本身就是 dict(上游已结构化),也可能是 str。
# 统一包装成与 _diagnose_exception 同构的字段,便于前端按一种模型渲染。
detail_payload: Any = exc.detail
error_detail: dict[str, Any]
if isinstance(detail_payload, dict) and "code" in detail_payload:
error_detail = detail_payload
else:
error_detail = {
"code": "http_exception",
"message": f"[{exc.status_code}] {detail_payload}",
"exception_type": "HTTPException",
"exception_repr": repr(exc),
"root_type": None,
"root_repr": None,
"hint": None,
"trace_tail": [],
"context": "intent_stream",
}
await record_foreground_diag(
request, stage="step1_intent", level="error", event="intent_http_error",
message=f"第一步意图生成失败:{error_detail.get('message')}",
detail=error_detail,
)
err = {
"status": "error",
"content": error_detail["message"],
"summary": None,
"error_detail": error_detail,
}
await _record_terminal(err)
yield f"data: {json.dumps(err, ensure_ascii=False)}\n\n"
except Exception as exc:
diagnostic = _diagnose_exception(exc, context="intent_stream")
logger.error("intent stream failed: %s", diagnostic, exc_info=exc)
await record_foreground_diag(
request, stage="step1_intent", level="error", event="intent_unexpected",
message=f"第一步意图生成未预期异常:{diagnostic.get('message') or exc}",
detail=diagnostic,
)
err = {
"status": "error",
"content": diagnostic["message"],
"summary": None,
"error_detail": diagnostic,
}
await _record_terminal(err)
yield f"data: {json.dumps(err, ensure_ascii=False)}\n\n"
return StreamingResponse(
stream_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
"Connection": "keep-alive",
},
)