FinDataPilot / app /agent /nodes /planner.py
Fin-DataPilot Deploy Bot
ci: da17dfe
ab26c67
Raw
History Blame Contribute Delete
14.4 kB
"""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,
}