Spaces:
Sleeping
Sleeping
File size: 14,416 Bytes
ab26c67 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 | """Planner node: pre-decompose the user question into a multi-step plan.
Pipeline:
planner → router → executor → reflector → (need_more)
↑ ↓
└─────────┘ advance through plan / replan if exhausted
↓
synthesizer
The planner LLM call sees the full question + the available skill list
and outputs a plan: a list of {goal, target_skill, args} steps to
execute in order. The skill router then walks the plan step by step
without re-asking the LLM, which is both faster and more coherent
than the original "decide next, execute, decide next" loop.
The reflector can still trigger a re-plan (by clearing the plan state)
when the current plan is exhausted and a follow-up is needed. This
combines the best of both worlds: explicit upfront planning + reactive
re-planning on unexpected outcomes.
"""
from __future__ import annotations
import json
import logging
import re
from typing import Any
from langchain_core.messages import HumanMessage, SystemMessage
from app.agent.state import AgentState
from app.config import get_settings
from app.llm import build_chat_model
from app.skills.registry import REGISTRY
logger = logging.getLogger(__name__)
PLANNER_PROMPT = """你是 Fin-DataPilot 的规划器(Planner)。基于用户的最新问题,**预先**把它拆成"一步步执行"的具体子任务,输出一个执行计划。
# 输入
- 用户的最新问题
- 可用 Skill 列表(带每个 skill 的参数 schema 摘要)
# 输出(严格 JSON,不带 markdown 代码块)
{{
"plan": [
{{"goal": "<这一步要达成什么目标>", "target_skill": "<skill 名或 null>", "args": {{...}}}}
],
"rationale": "<简短解释为什么这么拆>"
}}
# 规则
1. **简单问题**("茅台股价"、"今天的新闻")→ **1 步计划**
2. **复合问题**("涨停 + 市值最大 + 公告"、"宁德时代为什么跌 + 上下游")→ **2-5 步计划**
3. 每一步必须有具体的 `target_skill` + `args`(除非是"最后总结"步骤,target_skill 可以是 null)
4. **args 的语义占位(重要)**:当后面步骤要引用前一步的输出时,**必须**用占位符:
- `<step_0_top_stock>` — 第 0 步结果中按市值最大的那只股票的「名称 + 代码」组合
- `<step_0_top_name>` — 只取名称
- `<step_0_top_code>` — 只取代码
- `<step_0_first>` — 第 0 步的第一行(JSON 字符串)
示例(查"涨停 + 市值最大"那只的公告):
```
step 1: target="financial-query", args={"query": "今日A股涨停股票,按总市值降序排序", "limit": 5}
step 2: target="announcement-search", args={"query": "<step_0_top_stock> 最近公告", "days": 30, "limit": 10}
```
**绝对不要**自己脑补具体标的(如 "admin"、"unknown"、"待定")— 那样会查不到东西
5. **从问题里**提取关键限制(时间、范围、数量)— 用户给的时间窗口必须带进 args
6. **金融问题可以多 Skill 互补**:四个金融 Skill(financial-query / news-search / announcement-search / report-search)不是互斥关系。解释原因、判断影响、分析风险、查近期近况时,通常要先取结构化数据,再补新闻 / 公告 / 研报。
7. **anysearch 是低优先级兜底,不是禁用**:金融 Skill 优先;如果金融 Skill 返回为空、字段不全、没有覆盖用户问句,或需要公开网页/实时事实核查,可以在后续步骤使用 anysearch。
8. **不要** plan 一个"最后总结"步骤 — synthesizer 会自动整合
9. **不要**重复同一步的 args — 如果前一步已发起的 query 拿到 0 结果,应换一种自然问法、补其它金融 Skill,或最后用 anysearch 兜底,而不是用相同 args 再发一次
# 例子
用户:「涨停的股票中市值最大的那只最近的公告和研报」
输出:
{{
"plan": [
{{"goal": "找涨停且市值最大的股票", "target_skill": "financial-query",
"args": {{"query": "今日A股涨停股票,按总市值降序排序", "limit": "5"}}}},
{{"goal": "查那只股票的近期研报", "target_skill": "report-search",
"args": {{"query": "<step_0_top_name>的研报", "days": "30", "limit": "10"}}}},
{{"goal": "查那只股票的最近公告", "target_skill": "announcement-search",
"args": {{"query": "<step_0_top_name>的公告", "days": "30", "limit": "10"}}}}
],
"rationale": "先取 top 股票,再分别查它的研报和公告"
}}
用户:「茅台股价多少」
输出:
{{"plan": [{{"goal": "取最新价", "target_skill": "financial-query",
"args": {{"query": "贵州茅台 最新价"}}}}],
"rationale": "单步问题,一查即得"}}
用户:「宁德时代为什么跌」
输出:
{{
"plan": [
{{"goal": "取宁德时代近期行情和资金变化", "target_skill": "financial-query",
"args": {{"query": "宁德时代近期涨跌幅、成交额、换手率、主力资金净流入", "limit": "10"}}}},
{{"goal": "查宁德时代近期公告事件", "target_skill": "announcement-search",
"args": {{"query": "宁德时代近期公告", "days": "30", "limit": "10"}}}},
{{"goal": "查宁德时代近期新闻", "target_skill": "news-search",
"args": {{"query": "宁德时代为什么跌 近期新闻", "days": "30", "limit": "10"}}}},
{{"goal": "查宁德时代近期研报观点", "target_skill": "report-search",
"args": {{"query": "宁德时代近期研报观点", "days": "30", "limit": "10"}}}}
],
"rationale": "原因分析需要行情/资金、公告、新闻、研报互相印证"
}}
用户:「今天杭州天气怎么样」
输出:
{{"plan": [{{"goal": "查实时天气", "target_skill": "anysearch",
"args": {{"action": "search", "query": "杭州 今天天气", "max_results": 5}}}}],
"rationale": "实时问题,走联网搜索"}}
"""
def _try_parse_plan(text: str) -> dict[str, Any] | None:
"""Best-effort extraction of a plan JSON from the LLM output."""
text = text.strip()
# Strip markdown code fences if present
if "```" in text:
for fence in text.split("```"):
fence = fence.strip()
if fence.startswith("json"):
fence = fence[4:].strip()
if fence.startswith("{"):
text = fence
break
try:
obj = json.loads(text)
except json.JSONDecodeError:
# Try to find the first {...} block
m = re.search(r"\{.*\}", text, re.DOTALL)
if not m:
return None
try:
obj = json.loads(m.group(0))
except json.JSONDecodeError:
return None
if not isinstance(obj, dict):
return None
plan = obj.get("plan")
if not isinstance(plan, list):
return None
# Validate / coerce each step
clean: list[dict[str, Any]] = []
for step in plan:
if not isinstance(step, dict):
continue
clean.append({
"goal": str(step.get("goal", "")),
"target_skill": step.get("target_skill"),
"args": step.get("args", {}) if isinstance(step.get("args"), dict) else {},
})
return {"plan": clean, "rationale": obj.get("rationale", "")}
def _requests_both_announcement_and_report(user_query: str) -> bool:
"""True when the user explicitly asks for both announcements and reports."""
q = user_query or ""
has_announcement = any(k in q for k in ("公告", "披露", "announcement", "filing"))
has_report = any(k in q for k in ("研报", "研究报告", "research report"))
if not (has_announcement and has_report):
return False
# "公告或研报" means either source is acceptable; "公告和研报" means both.
if any(k in q for k in ("或", "或者", "二选一", "任一")):
return False
return True
def _normalize_plan_for_query(user_query: str, plan: list[dict[str, Any]]) -> list[dict[str, Any]]:
"""Patch common LLM under-planning for "top stock + 公告和研报".
The planner sometimes emits only announcement-search for a question
that asks for both announcements and reports. For this high-traffic
pattern, make the two dependent lookups explicit and deterministic.
"""
if not plan or not _requests_both_announcement_and_report(user_query):
return plan
try:
financial_idx = next(
i for i, step in enumerate(plan)
if step.get("target_skill") == "financial-query"
)
except StopIteration:
return plan
if financial_idx != 0:
# Placeholders are indexed by prior tool-call order. Keep this
# normalization conservative unless the financial query is first.
return plan
def _followup_step(skill: str) -> dict[str, Any]:
if skill == "report-search":
return {
"goal": "查询市值最大涨停股的近期研报",
"target_skill": "report-search",
"args": {"query": "<step_0_top_name>的研报", "days": "30", "limit": "10"},
}
return {
"goal": "查询市值最大涨停股的近期公告",
"target_skill": "announcement-search",
"args": {"query": "<step_0_top_name>的公告", "days": "30", "limit": "10"},
}
normalized: list[dict[str, Any]] = []
inserted = False
for idx, step in enumerate(plan):
if step.get("target_skill") in ("report-search", "announcement-search"):
continue
normalized.append(step)
if idx == financial_idx and not inserted:
normalized.append(_followup_step("report-search"))
normalized.append(_followup_step("announcement-search"))
inserted = True
return normalized
async def planner_node(state: AgentState) -> dict[str, Any]:
"""Decompose the user question into a multi-step plan.
Re-entry behavior: if the state already has a `plan` (from a prior
call or from a replan), do nothing. This is what makes replan work
— the reflector clears the plan to force a re-invocation.
"""
user_query = state.get("user_query", "")
history = state.get("history", []) or []
# If the planner is invoked a second time (replan), prepend the
# prior plan + all tool results so the LLM has the full context.
prior_plan = state.get("plan") or []
prior_calls = state.get("tool_calls") or []
if prior_plan:
# Replan: a previous plan exists but was exhausted. Re-decompose.
logger.info("planner: replanning (had %d-step plan)", len(prior_plan))
else:
logger.info("planner: first-time planning for query=%r", user_query[:80])
settings = get_settings()
llm = build_chat_model(settings, temperature=0.0)
history_text = "\n".join(f"[{m['role']}] {m['content']}" for m in history[-6:])
# Build a compact skill summary so the planner knows what's available.
skill_lines = []
for s in REGISTRY.list_specs():
params = ", ".join(
f"{p.name}{'' if p.required else '?'}: {p.type}" for p in s.parameters
)
skill_lines.append(f"- {s.name}({params}) — {s.description[:120]}")
skills_text = "\n".join(skill_lines) or "(无可用 Skill)"
# On replan, include prior steps + their results so the LLM can
# build a follow-up plan that picks up where we left off.
prior_text = ""
if prior_calls:
prior_text = "\n\n# 已完成的工具调用(按时间顺序)\n" + "\n".join(
f"### Step {i}: {c.get('name')}({json.dumps(c.get('args', {}), ensure_ascii=False)})\n"
f"Result summary: {json.dumps((c.get('result') or {}).get('data'), ensure_ascii=False)[:600]}"
for i, c in enumerate(prior_calls)
)
user_prompt = (
f"# 对话历史(最近 6 条)\n{history_text or '(无)'}\n\n"
f"# 用户最新问题\n{user_query}\n\n"
f"# 可用 Skill\n{skills_text}"
f"{prior_text}\n\n"
"请按 system prompt 中的契约输出 plan JSON。"
)
try:
resp = await llm.ainvoke(
[SystemMessage(content=PLANNER_PROMPT), HumanMessage(content=user_prompt)]
)
except Exception as exc: # noqa: BLE001
logger.exception("planner LLM call failed")
# Fallback: empty plan → router's LLM path will drive the
# question reactively. Don't emit a fake 1-step plan with
# null skill; that short-circuits the whole run.
return {
"plan": [],
"pending_step_index": 0,
"error": f"planner LLM call failed: {exc}",
}
content = resp.content if isinstance(resp.content, str) else str(resp.content)
parsed = _try_parse_plan(content)
if not parsed or not parsed.get("plan"):
logger.warning("planner: failed to parse plan (raw output: %r), falling back to reactive router", content[:300])
return {
"plan": [],
"pending_step_index": 0,
}
# Validate: every step's target_skill (if not None) must exist + be enabled.
clean_plan: list[dict[str, Any]] = []
for step in parsed["plan"]:
skill = step.get("target_skill")
if skill is None:
clean_plan.append(step)
continue
if not REGISTRY.get_spec(skill):
logger.warning("planner: unknown skill %r in plan, dropping step", skill)
continue
if not REGISTRY.is_enabled(skill):
logger.warning("planner: disabled skill %r in plan, dropping step", skill)
continue
clean_plan.append(step)
# Edge case: planner gave us an empty plan after validation —
# fall back to letting the router LLM handle it reactively.
if not clean_plan:
logger.warning("planner: every planned step was invalid, falling back to reactive router")
else:
clean_plan = _normalize_plan_for_query(user_query, clean_plan)
logger.info(
"planner: produced %d-step plan: %s",
len(clean_plan),
[s.get("target_skill") for s in clean_plan],
)
return {
"plan": clean_plan,
"pending_step_index": 0,
# Clear any stale hint from a prior reflector turn.
"next_skill_hint": None,
"next_args_hint": None,
}
|