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,
    }