Files
Cloud-Tour-to-Libo/app/graph_qa_engine.py
T

2848 lines
118 KiB
Python
Raw 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.
"""LLM-to-Cypher graph QA engine for enterprise knowledge-base questions."""
from __future__ import annotations
import asyncio
import copy
import json
import re
import time
from typing import Any
from falkordb import FalkorDB
from app.config import settings
from app.db import get_agent_settings
from app.llm_client import LlmClient
GRAPH_QA_MAX_LIMIT = 300
SCHEMA_CACHE_TTL_SECONDS = 300
QA_RESPONSE_CACHE_TTL_SECONDS = 600
CYPHER_CACHE_TTL_SECONDS = 1800
GRAPH_QA_PLANNER_TIMEOUT_SECONDS = 2.5
GRAPH_QA_CYPHER_TIMEOUT_SECONDS = 4.5
GRAPH_QA_REPAIR_TIMEOUT_SECONDS = 5.0
GRAPH_QA_ANSWER_TIMEOUT_SECONDS = 6.0
_SCHEMA_CACHE: dict[str, tuple[float, dict[str, Any]]] = {}
_QA_RESPONSE_CACHE: dict[str, tuple[float, dict[str, Any]]] = {}
_CYPHER_CACHE: dict[str, tuple[float, dict[str, Any]]] = {}
WRITE_KEYWORDS = re.compile(
r"\b(CREATE|MERGE|SET|DELETE|DETACH|DROP|REMOVE|CALL\s+DB\.IDX|CALL\s+DB\.CONSTRAINT)\b",
re.IGNORECASE,
)
READ_START = re.compile(r"^(MATCH|RETURN|CALL)\b", re.IGNORECASE)
GRAPH_QA_PLANNER_SYS = """你是百姓惠知识图谱客服问答的路由规划器。只输出 JSON:
{"intent":"...","route":"...","confidence":0.0,"reason":"...","need_llm_graph":false}
任务:判断用户问题应该先走哪条链路,不生成答案、不生成 Cypher。
可选 intent:
route_catalog、price_quote、fee_detail、vehicle_catalog、route_multihop_detail、
condition_route_advice、route_compare_suitability、route_compare_scenic_count、
route_context_advice、general_graph_qa。
可选 route:
template_graph_query、deterministic_context_llm_answer、llm_graph_qa。
原则:高频客服问题优先 template_graph_query;需要推荐/解释但图谱可固定取证时选 deterministic_context_llm_answer;
无法用模板覆盖的长尾问题才选 llm_graph_qa。不要臆造线路名。"""
GRAPH_QA_CYPHER_SYS = """你是 FalkorDB/RedisGraph 只读 Cypher 规划器。只输出 JSON:
{"cypher":"...","reason":"...","answer_focus":"..."}
规则:只允许 MATCH/RETURN/CALL;禁止 CREATE/MERGE/SET/DELETE/DETACH/DROP/REMOVE;必须 LIMIT;中文用 CONTAINS。
只给一条最短可执行查询,cypher 控制在 800 字以内,不要输出示例、解释或多条备选。
百姓惠常用标签:TourProduct、ProductDay、RouteStop、ScenicAttraction、TravelItem、ProductPricePlan、PolicyRule、SalesScript。
价格:TourProduct-[:PRODUCT_HAS_PRICE_PLAN]->ProductPricePlan,字段 adult_price_min/max、child_price_min/max、single_room_diff_min/max。
费用/小交通:TourProduct-[:PRODUCT_USES_TRAVEL_ITEM]->TravelItem;TravelItem 字段 name、type、subtype、category、price、adult_price、unit、raw_evidence、requires_supplier_confirm。
黄小西/小西:按黄果树、小七孔、西江相关线路匹配。"""
GRAPH_QA_REPAIR_SYS = """你是 FalkorDB/RedisGraph Cypher 修复器。只输出 JSON。
根据执行错误修复上一条只读 Cypher。仍然必须只读、以 MATCH/RETURN/CALL 开头、带 LIMIT。
字段:cypher、reason、answer_focus。"""
GRAPH_QA_ANSWER_SYS = """你是百姓惠客服图谱问答助手。只输出 JSON。
基于 graph_result 回答,不编造价格、余位、房型、车辆或承诺;证据不足就提示二次核实。
customer_reply 控制在 100 字内,answer 控制在 360 字内。
字段:answer、customer_reply、confidence、follow_up_questions、risk_notes、evidence_cards。"""
def _get_graph(graph_name: str):
db = FalkorDB(host=settings.falkordb_host, port=settings.falkordb_port)
return db.select_graph(graph_name)
def _safe_int(value: Any, default: int) -> int:
try:
return int(value)
except Exception:
return default
async def _chat_json_with_retry(
client: LlmClient,
system_prompt: str,
user_prompt: str,
*,
attempts: int = 2,
) -> dict[str, Any]:
last_exc: Exception | None = None
for attempt in range(max(1, attempts)):
try:
return await asyncio.to_thread(client.chat_json, system_prompt, user_prompt)
except Exception as exc: # noqa: BLE001
last_exc = exc
if attempt >= attempts - 1:
break
await asyncio.sleep(0.8 * (attempt + 1))
raise last_exc or RuntimeError("LLM JSON 调用失败")
async def _chat_json_timed(
client: LlmClient,
system_prompt: str,
user_prompt: str,
*,
timeout_seconds: float,
attempts: int = 1,
) -> dict[str, Any]:
return await asyncio.wait_for(
_chat_json_with_retry(client, system_prompt, user_prompt, attempts=attempts),
timeout=timeout_seconds,
)
async def _graph_qa_client(max_tokens: int = 1600, *, use_config_max_tokens: bool = True) -> LlmClient | None:
try:
cfg = await get_agent_settings()
except Exception:
cfg = {}
api_access = cfg.get("api_access") if isinstance(cfg, dict) else {}
qa_llm = (api_access or {}).get("qa_llm") if isinstance(api_access, dict) else {}
if isinstance(qa_llm, dict) and qa_llm.get("base_url") and qa_llm.get("api_key"):
configured_max_tokens = _safe_int(qa_llm.get("max_tokens"), max_tokens)
return LlmClient(
qa_llm["base_url"],
qa_llm["api_key"],
qa_llm.get("model") or settings.llm_model or "deepseek-chat",
timeout=_safe_int(qa_llm.get("timeout"), settings.llm_timeout_seconds or 45),
max_tokens=min(configured_max_tokens, max_tokens) if use_config_max_tokens else max_tokens,
)
global_cfg = cfg.get("global") if isinstance(cfg, dict) else {}
if isinstance(global_cfg, dict) and global_cfg.get("base_url") and global_cfg.get("api_key"):
return LlmClient(
global_cfg["base_url"],
global_cfg["api_key"],
global_cfg.get("model") or settings.llm_model or "deepseek-chat",
timeout=_safe_int(global_cfg.get("timeout"), settings.llm_timeout_seconds or 45),
max_tokens=max_tokens,
)
if settings.llm_api_base and settings.llm_api_key:
return LlmClient(
settings.llm_api_base,
settings.llm_api_key,
settings.llm_model or "deepseek-chat",
timeout=settings.llm_timeout_seconds,
max_tokens=max_tokens,
)
extract = cfg.get("extract") if isinstance(cfg, dict) else {}
extract = extract if isinstance(extract, dict) else {}
models = extract.get("models") or {}
models = models if isinstance(models, dict) else {}
keys = [extract.get("aggregator")] + [
key for key, value in models.items() if isinstance(value, dict) and value.get("enabled")
] + list(models.keys())
for key in keys:
model_cfg = models.get(key or "")
if isinstance(model_cfg, dict) and model_cfg.get("base_url") and model_cfg.get("api_key"):
return LlmClient(
model_cfg["base_url"],
model_cfg["api_key"],
model_cfg.get("model") or settings.llm_model or "deepseek-chat",
timeout=int(extract.get("timeout") or settings.llm_timeout_seconds or 60),
max_tokens=max_tokens,
)
return None
def _assert_read_only(cypher: str) -> None:
if not cypher:
raise ValueError("Cypher query required")
if WRITE_KEYWORDS.search(cypher):
raise ValueError("Read-only queries only")
if not READ_START.match(cypher):
raise ValueError("Only MATCH, RETURN, or CALL queries allowed")
def _clean_cypher(cypher: Any, limit: int) -> str:
text = str(cypher or "").strip()
text = re.sub(r"^```(?:cypher)?\s*", "", text, flags=re.I).strip()
text = text.removesuffix("```").strip().rstrip(";")
if not re.search(r"\bLIMIT\s+\d+\b", text, flags=re.IGNORECASE):
text = f"{text} LIMIT {limit}"
_assert_read_only(text)
return text
def _compact_scalar(value: Any, limit: int = 500) -> Any:
if value is None or isinstance(value, (int, float, bool)):
return value
if isinstance(value, (list, tuple)):
return [_compact_scalar(item, 180) for item in value[:20]]
text = re.sub(r"\s+", " ", str(value)).strip()
return text[:limit]
def _node_id(node: Any) -> str:
value = getattr(node, "id", None)
return str(value if value is not None else node)
def _node_summary(node: Any) -> dict[str, Any]:
props = dict(getattr(node, "properties", None) or {})
labels = list(getattr(node, "labels", None) or [])
title = ""
for key in (
"display_name", "name", "title", "label", "product_name",
"route_line_name", "natural_key", "product_id", "item_id",
):
if props.get(key):
title = str(props[key])
break
if not title:
title = f"{labels[0] if labels else 'Node'} #{_node_id(node)}"
return {
"id": _node_id(node),
"labels": labels,
"title": title[:160],
"properties": {str(k): _compact_scalar(v, 700) for k, v in props.items()},
}
def _edge_summary(edge: Any) -> dict[str, Any]:
props = dict(getattr(edge, "properties", None) or {})
return {
"id": str(getattr(edge, "id", "")),
"type": str(getattr(edge, "relation", "") or ""),
"from": _node_id(getattr(edge, "src_node", "")),
"to": _node_id(getattr(edge, "dest_node", "")),
"properties": {str(k): _compact_scalar(v, 300) for k, v in props.items()},
}
def _compact_value(value: Any) -> Any:
if hasattr(value, "labels") and hasattr(value, "properties"):
return {"node": _node_summary(value)}
if hasattr(value, "src_node") and hasattr(value, "dest_node"):
return {"relationship": _edge_summary(value)}
if (
hasattr(value, "nodes")
and callable(getattr(value, "nodes"))
and hasattr(value, "edges")
and callable(getattr(value, "edges"))
):
return {
"path": {
"nodes": [_node_summary(node) for node in value.nodes()],
"relationships": [_edge_summary(edge) for edge in value.edges()],
}
}
if isinstance(value, (list, tuple)):
return [_compact_value(item) for item in value[:20]]
if isinstance(value, dict):
return {str(k): _compact_value(v) for k, v in list(value.items())[:30]}
return _compact_scalar(value)
def _extract_graph_items(rows: list[list[Any]]) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]:
nodes: dict[str, dict[str, Any]] = {}
edges: dict[str, dict[str, Any]] = {}
def walk(value: Any) -> None:
if hasattr(value, "labels") and hasattr(value, "properties"):
nodes.setdefault(_node_id(value), _node_summary(value))
return
if hasattr(value, "src_node") and hasattr(value, "dest_node"):
edge = _edge_summary(value)
key = edge["id"] or f"{edge['from']}-{edge['type']}-{edge['to']}"
edges.setdefault(key, edge)
return
if (
hasattr(value, "nodes")
and callable(getattr(value, "nodes"))
and hasattr(value, "edges")
and callable(getattr(value, "edges"))
):
for node in value.nodes():
walk(node)
for edge in value.edges():
walk(edge)
return
if isinstance(value, (list, tuple)):
for item in value:
walk(item)
for row in rows:
walk(row)
return list(nodes.values()), list(edges.values())
def _header_names(header: Any) -> list[str]:
out: list[str] = []
for item in header or []:
if isinstance(item, (list, tuple)) and len(item) > 1:
out.append(str(item[1]))
else:
out.append(str(item))
return out
def _run_cypher(graph_name: str, cypher: str, limit: int) -> dict[str, Any]:
graph = _get_graph(graph_name)
result = graph.query(cypher)
rows = list(result.result_set)[:limit]
nodes, relationships = _extract_graph_items(rows)
compact_rows = [[_compact_value(value) for value in row] for row in rows[:40]]
return {
"cypher": cypher,
"columns": _header_names(result.header),
"rows": compact_rows,
"row_count": len(rows),
"nodes": nodes,
"relationships": relationships,
}
FEE_QUESTION_TERMS = (
"费用", "自费", "自理", "另付", "不含", "包含", "小交通", "观光车", "环保车",
"电瓶车", "索道", "扶梯", "游船", "门票", "二消", "必付", "可选",
)
PRICE_QUESTION_TERMS = (
"多少钱", "价格", "报价", "成人价", "儿童价", "小孩价", "房差", "单房差", "结算价",
)
VEHICLE_QUESTION_TERMS = (
"车型", "车辆", "用车", "车有哪些", "什么车", "商务车", "大巴", "中巴", "小包团用车",
)
COMPLEX_QUESTION_TERMS = (
"推荐", "对比", "比较", "为什么", "怎么安排", "如何安排", "适合", "帮我", "方案",
"预算", "老人", "亲子", "团队", "定制", "如果", "同时", "综合", "详细行程",
)
FAMILY_CONDITION_TERMS = ("老人", "老年", "长辈", "小孩", "孩子", "儿童", "亲子", "家庭", "一家")
COMFORT_CONDITION_TERMS = ("不累", "不要太累", "太累", "轻松", "舒适", "少走路", "慢一点", "不赶", "休闲", "轻奢")
BUDGET_CONDITION_TERMS = ("预算", "便宜", "划算", "性价比", "不贵", "低价", "经济")
ROUTE_COMPARE_TERMS = ("哪个", "哪条", "对比", "比较", "更多", "更少", "包含")
SCENIC_QUESTION_TERMS = (
"景区", "景点", "去哪", "去哪里", "去哪些", "游玩", "打卡", "途经", "经过",
)
HOTEL_QUESTION_TERMS = (
"酒店", "住宿", "入住", "住哪里", "住哪", "附近酒店", "附近入住", "附近住宿",
)
SCENIC_HINT_ALIASES: dict[str, tuple[str, ...]] = {
"荔波小七孔景区": ("小七孔", "荔波"),
"西江千户苗寨景区": ("西江", "千户苗寨", "苗寨"),
"黄果树旅游景区": ("黄果树", "瀑布"),
"梵净山景区": ("梵净山", "镇梵"),
"青岩古镇": ("青岩", "青岩古镇"),
"镇远古城": ("镇远", "镇远古城"),
"贵阳": ("贵阳",),
"安顺": ("安顺",),
}
FEE_EVIDENCE_TERMS = (
"raw_evidence", "ScenicTransport", "ScenicOptional", "景区费用项目", "观光车",
"环保车", "电瓶车", "小交通", "索道", "扶梯", "游船", "门票", "自费", "不含", "另付",
)
def _is_fee_question(question: str) -> bool:
return any(term in question for term in FEE_QUESTION_TERMS)
def _is_price_question(question: str) -> bool:
return any(term in question for term in PRICE_QUESTION_TERMS)
def _is_vehicle_question(question: str) -> bool:
return any(term in question for term in VEHICLE_QUESTION_TERMS)
def _has_duration_question(question: str) -> bool:
return bool(_route_duration_clauses(question)) or any(term in question for term in ("几天", "多少天", "天数", "玩几天", "行程几天"))
def _is_scenic_question(question: str) -> bool:
return any(term in question for term in SCENIC_QUESTION_TERMS)
def _is_hotel_question(question: str) -> bool:
return any(term in question for term in HOTEL_QUESTION_TERMS)
def _is_route_multihop_question(question: str) -> bool:
if not _route_name_terms(question) or not _is_hotel_question(question):
return False
facets = [
_is_price_question(question),
_has_duration_question(question),
_is_scenic_question(question),
_is_hotel_question(question),
]
return sum(1 for item in facets if item) >= 3
def _is_route_compare_scenic_count_question(question: str) -> bool:
terms = _route_name_terms(question)
return (
len(terms) >= 2
and any(term in question for term in ROUTE_COMPARE_TERMS)
and any(term in question for term in ("景点", "景区", "包含", "去的地方", "玩的地方"))
and any(term in question for term in ("更多", "多", "少", "哪个", "哪条"))
)
def _is_route_compare_suitability_question(question: str) -> bool:
terms = _route_name_terms(question)
slots = _condition_slots(question)
return (
len(terms) >= 2
and (slots["family"] or slots["comfort"] or slots["budget"])
and any(term in question for term in ("哪个", "哪条", "对比", "比较", "适合", "推荐", "为什么"))
)
def _condition_slots(question: str) -> dict[str, bool]:
return {
"family": any(term in question for term in FAMILY_CONDITION_TERMS),
"comfort": any(term in question for term in COMFORT_CONDITION_TERMS),
"budget": any(term in question for term in BUDGET_CONDITION_TERMS),
"recommend": any(term in question for term in ("推荐", "适合", "有没有", "哪条", "帮我", "怎么选")),
}
def _is_condition_route_advice_question(question: str) -> bool:
slots = _condition_slots(question)
return (
not _route_name_terms(question)
and slots["recommend"]
and (slots["family"] or slots["comfort"] or slots["budget"])
)
def _normalized_question(question: str) -> str:
text = re.sub(r"\s+", "", question.strip().lower())
return (
text.replace("路线", "线路")
.replace("产品", "线路")
.replace("报价", "价格")
.replace("成人价", "价格")
.replace("小孩价", "儿童价")
)
PLANNER_ROUTE_BY_INTENT = {
"route_catalog": "template_graph_query",
"price_quote": "template_graph_query",
"fee_detail": "template_graph_query",
"vehicle_catalog": "template_graph_query",
"route_multihop_detail": "template_graph_query",
"condition_route_advice": "template_graph_query",
"route_compare_suitability": "template_graph_query",
"route_compare_scenic_count": "template_graph_query",
"route_context_advice": "deterministic_context_llm_answer",
"general_graph_qa": "llm_graph_qa",
}
def _planner_enabled_from_context(customer_context: dict[str, Any] | None) -> bool:
if not isinstance(customer_context, dict):
return False
value = customer_context.get("llm_planner")
if value is None:
value = customer_context.get("force_llm_planner")
return str(value).strip().lower() in {"1", "true", "yes", "on", "启用", "强制"}
def _should_use_llm_planner(intent: dict[str, Any], customer_context: dict[str, Any] | None) -> bool:
if _planner_enabled_from_context(customer_context):
return True
route = str(intent.get("route") or "")
confidence = _safe_confidence(intent.get("confidence"), 0.0)
return route == "llm_graph_qa" or confidence < 0.7
def _planner_trace(
*,
used: bool,
skipped_reason: str = "",
latency_ms: int = 0,
error: str = "",
decision: dict[str, Any] | None = None,
) -> dict[str, Any]:
return {
"used": used,
"latency_ms": latency_ms,
"timeout_ms": round(GRAPH_QA_PLANNER_TIMEOUT_SECONDS * 1000),
"skipped_reason": skipped_reason,
"error": error,
"decision": decision or {},
}
def _normalize_planner_decision(base_intent: dict[str, Any], decision: dict[str, Any]) -> dict[str, Any]:
planned_intent = str(decision.get("intent") or "").strip()
if planned_intent not in PLANNER_ROUTE_BY_INTENT:
return base_intent
confidence = _safe_confidence(decision.get("confidence"), _safe_confidence(base_intent.get("confidence"), 0.55))
route = PLANNER_ROUTE_BY_INTENT[planned_intent]
return {
**base_intent,
"intent": planned_intent,
"route": route,
"confidence": confidence,
"complexity": str(decision.get("complexity") or base_intent.get("complexity") or "planned"),
"reason": str(decision.get("reason") or base_intent.get("reason") or "LLM 规划器选择链路。")[:260],
"planner_need_llm_graph": bool(decision.get("need_llm_graph")) or route == "llm_graph_qa",
}
async def _llm_plan_intent(
question: str,
graph_name: str,
base_intent: dict[str, Any],
customer_context: dict[str, Any] | None,
) -> tuple[dict[str, Any], dict[str, Any]]:
if not _should_use_llm_planner(base_intent, customer_context):
return base_intent, _planner_trace(used=False, skipped_reason="high_confidence_template_or_context_route")
client = await _graph_qa_client(max_tokens=320, use_config_max_tokens=False)
if client is None:
return base_intent, _planner_trace(used=False, skipped_reason="llm_not_configured")
client.timeout = min(float(getattr(client, "timeout", GRAPH_QA_PLANNER_TIMEOUT_SECONDS) or GRAPH_QA_PLANNER_TIMEOUT_SECONDS), GRAPH_QA_PLANNER_TIMEOUT_SECONDS)
payload = {
"question": question,
"graph_name": graph_name,
"rule_intent": base_intent,
"route_terms": _route_name_terms(question),
"condition_slots": _condition_slots(question),
"allowed_intents": list(PLANNER_ROUTE_BY_INTENT.keys()),
"customer_context": customer_context or {},
}
started = time.perf_counter()
try:
decision = await _chat_json_timed(
client,
GRAPH_QA_PLANNER_SYS,
json.dumps(payload, ensure_ascii=False),
timeout_seconds=GRAPH_QA_PLANNER_TIMEOUT_SECONDS,
attempts=1,
)
except asyncio.TimeoutError:
latency_ms = max(1, round((time.perf_counter() - started) * 1000))
return base_intent, _planner_trace(used=True, latency_ms=latency_ms, error="planner_timeout")
except Exception as exc: # noqa: BLE001
latency_ms = max(1, round((time.perf_counter() - started) * 1000))
return base_intent, _planner_trace(used=True, latency_ms=latency_ms, error=str(exc)[:220])
latency_ms = max(1, round((time.perf_counter() - started) * 1000))
planned = _normalize_planner_decision(base_intent, decision)
return planned, _planner_trace(used=True, latency_ms=latency_ms, decision=decision)
def _classify_graph_qa_intent(question: str) -> dict[str, Any]:
normalized = _normalized_question(question)
terms = _route_name_terms(question)
duration_clauses = _route_duration_clauses(question)
complex_hint = len(question) > 72 or any(term in question for term in COMPLEX_QUESTION_TERMS)
intent = {
"intent": "general_graph_qa",
"route": "llm_graph_qa",
"confidence": 0.48 if complex_hint else 0.58,
"complexity": "complex" if complex_hint else "open",
"slots": {
"route_terms": terms,
"has_duration": bool(duration_clauses),
},
"reason": "未命中标准模板,使用 LLM 生成 Cypher 做通用图查询。",
}
if _is_route_list_question(question):
return {
**intent,
"intent": "route_catalog",
"route": "template_graph_query",
"confidence": 0.95,
"complexity": "standard",
"reason": "线路清单类问题可由 TourProduct 固定模板回答。",
}
if _is_condition_route_advice_question(question):
return {
**intent,
"intent": "condition_route_advice",
"route": "template_graph_query",
"confidence": 0.82,
"complexity": "standard_condition_filter",
"slots": {
**intent["slots"],
"condition_slots": _condition_slots(question),
},
"reason": "未指定线路但给出家庭/舒适/预算条件,使用固定图查询筛选候选线路并模板回答。",
}
if _is_route_compare_suitability_question(question):
return {
**intent,
"intent": "route_compare_suitability",
"route": "template_graph_query",
"confidence": 0.84,
"complexity": "standard_suitability_compare",
"slots": {
**intent["slots"],
"condition_slots": _condition_slots(question),
},
"reason": "多线路适配度对比已识别条件词,可固定取证后用评分模板快速回答。",
}
if _is_route_compare_scenic_count_question(question):
return {
**intent,
"intent": "route_compare_scenic_count",
"route": "template_graph_query",
"confidence": 0.86,
"complexity": "standard_compare",
"reason": "线路景点数量对比可由固定取证统计 ScenicAttraction 数量回答。",
}
if _is_vehicle_question(question) and any(term in normalized for term in ("车型", "车辆", "用车", "车")):
return {
**intent,
"intent": "vehicle_catalog",
"route": "template_graph_query",
"confidence": 0.88,
"complexity": "standard",
"reason": "车辆/车型类问题可由 TravelItem 用车项目固定模板回答。",
}
if _is_route_multihop_question(question):
return {
**intent,
"intent": "route_multihop_detail",
"route": "template_graph_query",
"confidence": 0.91,
"complexity": "standard_multihop",
"reason": "线路价格、天数、途经景区和附近酒店属于标准多跳图谱查询,可用固定模板快速回答。",
}
if _is_price_question(question) and terms:
return {
**intent,
"intent": "price_quote",
"route": "template_graph_query",
"confidence": 0.9,
"complexity": "standard",
"reason": "线路报价类问题已识别线路实体和天数,可直接查 ProductPricePlan。",
}
if _is_fee_question(question) and terms:
return {
**intent,
"intent": "fee_detail",
"route": "template_graph_query",
"confidence": 0.86,
"complexity": "standard",
"reason": "费用/小交通类问题已识别线路实体,可直接查 TravelItem。",
}
if complex_hint and terms:
return {
**intent,
"intent": "route_context_advice",
"route": "deterministic_context_llm_answer",
"confidence": 0.76,
"complexity": "complex_with_entities",
"reason": "复杂推荐/对比问题已识别线路实体,先用固定 Cypher 取证据,再让 LLM 基于证据组织答案。",
}
return intent
def _has_fee_evidence(graph_result: dict[str, Any]) -> bool:
if not graph_result.get("row_count"):
return False
text = json.dumps(
{
"rows": graph_result.get("rows") or [],
"nodes": graph_result.get("nodes") or [],
"relationships": graph_result.get("relationships") or [],
},
ensure_ascii=False,
default=str,
)
return any(term in text for term in FEE_EVIDENCE_TERMS)
def _cypher_quote(value: str) -> str:
return value.replace("\\", "\\\\").replace("'", "\\'")
def _is_valid_route_term(term: str) -> bool:
text = term.strip()
if len(text) < 2:
return False
noisy_fragments = (
"适合", "推荐", "有没有", "不要", "不累", "太累", "轻松", "舒适", "老人", "小孩",
"孩子", "儿童", "家庭", "预算", "便宜", "划算", "哪些", "什么", "哪个", "多少",
"可以", "附近", "入住", "酒店",
)
if any(fragment in text for fragment in noisy_fragments):
return False
return True
def _route_name_terms(question: str) -> list[str]:
terms: list[str] = []
if "小西镇梵" in question:
terms.append("小西镇梵")
if "黄小西镇梵" in question:
terms.append("黄小西镇梵")
if "黄小西" in question:
terms.append("黄小西")
elif "小西" in question:
terms.append("小西")
for term in ("黄果树", "小七孔", "荔波", "西江", "千户苗寨", "青岩", "梵净山", "镇远"):
if term in question:
terms.append(term)
if not terms:
for match in re.findall(r"[\u4e00-\u9fa5A-Za-z0-9]{2,12}(?:游|线路|产品|团)", question):
terms.append(match.replace("线路", "").replace("产品", "").replace("团", ""))
out: list[str] = []
for term in terms:
if _is_valid_route_term(term) and term not in out:
out.append(term)
return out[:6]
def _route_duration_clauses(question: str) -> list[str]:
match = re.search(r"(\d+)\s*(?:日|天)", question)
days = int(match.group(1)) if match else None
chinese_days = {"一": 1, "二": 2, "两": 2, "三": 3, "四": 4, "五": 5, "六": 6, "七": 7}
for text, value in chinese_days.items():
if f"{text}日" in question or f"{text}天" in question:
days = value
break
if not days:
return []
return [
f"p.duration_days = {days}",
f"p.name CONTAINS '{days}日'",
f"p.name CONTAINS '{days}天'",
f"p.name CONTAINS '{list(chinese_days.keys())[list(chinese_days.values()).index(days)]}日'" if days in chinese_days.values() else "",
]
def _fee_fallback_cypher(question: str, limit: int) -> str | None:
terms = _route_name_terms(question)
if not terms:
return None
name_clauses: list[str] = []
for term in terms:
quoted = _cypher_quote(term)
name_clauses.extend([
f"p.name CONTAINS '{quoted}'",
f"p.display_name CONTAINS '{quoted}'",
f"p.product_name CONTAINS '{quoted}'",
])
duration_clauses = [item for item in _route_duration_clauses(question) if item]
where_parts = [f"({' OR '.join(name_clauses)})"]
if duration_clauses:
where_parts.append(f"({' OR '.join(duration_clauses)})")
where_clause = " AND ".join(where_parts)
return f"""
MATCH (p:TourProduct)
WHERE {where_clause}
OPTIONAL MATCH (p)-[:PRODUCT_USES_TRAVEL_ITEM]->(direct_item:TravelItem)
OPTIONAL MATCH (p)-[:PRODUCT_HAS_DAY]->(d:ProductDay)-[:DAY_HAS_STOP]->(s:RouteStop)
OPTIONAL MATCH (s)-[:STOP_VISITS_ATTRACTION]->(a:ScenicAttraction)
OPTIONAL MATCH (s)-[:STOP_USES_TRAVEL_ITEM]->(stop_item:TravelItem)
OPTIONAL MATCH (a)-[:ATTRACTION_HAS_ITEM]->(attraction_item:TravelItem)
RETURN
p.name AS product_name,
p.product_id AS product_id,
p.duration_days AS duration_days,
collect(DISTINCT {{
name: direct_item.name,
type: direct_item.type,
subtype: direct_item.subtype,
category: direct_item.category,
status: direct_item.default_status,
price: direct_item.price,
adult_price: direct_item.adult_price,
unit: direct_item.unit,
raw_evidence: direct_item.raw_evidence,
requires_supplier_confirm: direct_item.requires_supplier_confirm
}}) AS product_items,
collect(DISTINCT {{
stop_name: s.name,
attraction_name: a.name,
name: stop_item.name,
type: stop_item.type,
subtype: stop_item.subtype,
category: stop_item.category,
status: stop_item.default_status,
price: stop_item.price,
adult_price: stop_item.adult_price,
unit: stop_item.unit,
raw_evidence: stop_item.raw_evidence,
requires_supplier_confirm: stop_item.requires_supplier_confirm
}}) AS stop_items,
collect(DISTINCT {{
stop_name: s.name,
attraction_name: a.name,
name: attraction_item.name,
type: attraction_item.type,
subtype: attraction_item.subtype,
category: attraction_item.category,
status: attraction_item.default_status,
price: attraction_item.price,
adult_price: attraction_item.adult_price,
unit: attraction_item.unit,
raw_evidence: attraction_item.raw_evidence,
requires_supplier_confirm: attraction_item.requires_supplier_confirm
}}) AS attraction_items
LIMIT {limit}
""".strip()
def _price_fallback_cypher(question: str, limit: int) -> str | None:
terms = _route_name_terms(question)
if not terms:
return None
name_clauses: list[str] = []
for term in terms:
quoted = _cypher_quote(term)
name_clauses.extend([
f"p.name CONTAINS '{quoted}'",
f"p.display_name CONTAINS '{quoted}'",
f"p.product_name CONTAINS '{quoted}'",
])
duration_clauses = [item for item in _route_duration_clauses(question) if item]
where_parts = [f"({' OR '.join(name_clauses)})"]
if duration_clauses:
where_parts.append(f"({' OR '.join(duration_clauses)})")
where_clause = " AND ".join(where_parts)
return f"""
MATCH (p:TourProduct)
WHERE {where_clause}
OPTIONAL MATCH (p)-[:PRODUCT_HAS_PRICE_PLAN]->(plan:ProductPricePlan)
RETURN
p.name AS product_name,
p.product_id AS product_id,
p.duration_days AS duration_days,
collect(DISTINCT {{
plan_name: plan.plan_name,
display_name: plan.display_name,
adult_price_min: plan.adult_price_min,
adult_price_max: plan.adult_price_max,
child_price_min: plan.child_price_min,
child_price_max: plan.child_price_max,
single_room_diff_min: plan.single_room_diff_min,
single_room_diff_max: plan.single_room_diff_max,
adult_settlement_source: plan.adult_settlement_source,
child_settlement_source: plan.child_settlement_source
}}) AS price_plans
LIMIT {limit}
""".strip()
def _route_catalog_cypher(limit: int) -> str:
return f"""
MATCH (p:TourProduct)
RETURN
p.name AS product_name,
p.display_name AS display_name,
p.product_id AS product_id,
p.duration_days AS duration_days
ORDER BY p.duration_days, p.name
LIMIT {limit}
""".strip()
def _vehicle_catalog_cypher(limit: int) -> str:
return f"""
MATCH (item:TravelItem)
WHERE item.name CONTAINS '车'
OR item.subtype CONTAINS '用车'
OR item.category CONTAINS '用车'
RETURN
item.name AS name,
item.type AS type,
item.subtype AS subtype,
item.price AS price,
item.unit AS unit,
item.raw_evidence AS raw_evidence
ORDER BY item.name
LIMIT {limit}
""".strip()
def _route_context_cypher(question: str, limit: int) -> str | None:
terms = _route_name_terms(question)
if not terms:
return None
name_clauses: list[str] = []
for term in terms:
quoted = _cypher_quote(term)
name_clauses.extend([
f"p.name CONTAINS '{quoted}'",
f"p.display_name CONTAINS '{quoted}'",
f"p.product_name CONTAINS '{quoted}'",
])
where_clause = f"({' OR '.join(name_clauses)})"
return f"""
MATCH (p:TourProduct)
WHERE {where_clause}
OPTIONAL MATCH (p)-[:PRODUCT_HAS_PRICE_PLAN]->(plan:ProductPricePlan)
OPTIONAL MATCH (p)-[:PRODUCT_HAS_DAY]->(d:ProductDay)-[:DAY_HAS_STOP]->(s:RouteStop)
OPTIONAL MATCH (s)-[:STOP_VISITS_ATTRACTION]->(a:ScenicAttraction)
RETURN
p.name AS product_name,
p.product_id AS product_id,
p.duration_days AS duration_days,
collect(DISTINCT {{
plan_name: plan.plan_name,
adult_price_min: plan.adult_price_min,
adult_price_max: plan.adult_price_max,
child_price_min: plan.child_price_min,
child_price_max: plan.child_price_max,
single_room_diff_min: plan.single_room_diff_min,
single_room_diff_max: plan.single_room_diff_max
}}) AS price_plans,
collect(DISTINCT {{
day: d.day_index,
stop_name: s.name,
attraction_name: a.name,
sequence: s.sequence
}}) AS route_points,
count(DISTINCT a) AS attraction_count
LIMIT {min(limit, 40)}
""".strip()
def _condition_context_cypher(limit: int) -> str:
return f"""
MATCH (p:TourProduct)
OPTIONAL MATCH (p)-[:PRODUCT_HAS_PRICE_PLAN]->(plan:ProductPricePlan)
OPTIONAL MATCH (p)-[:PRODUCT_HAS_DAY]->(d:ProductDay)-[:DAY_HAS_STOP]->(s:RouteStop)
OPTIONAL MATCH (s)-[:STOP_VISITS_ATTRACTION]->(a:ScenicAttraction)
RETURN
p.name AS product_name,
p.product_id AS product_id,
p.duration_days AS duration_days,
collect(DISTINCT {{
plan_name: plan.plan_name,
adult_price_min: plan.adult_price_min,
adult_price_max: plan.adult_price_max,
child_price_min: plan.child_price_min,
child_price_max: plan.child_price_max,
single_room_diff_min: plan.single_room_diff_min,
single_room_diff_max: plan.single_room_diff_max
}}) AS price_plans,
collect(DISTINCT {{
day: d.day_index,
stop_name: s.name,
attraction_name: a.name,
sequence: s.sequence
}}) AS route_points,
count(DISTINCT a) AS attraction_count
ORDER BY p.duration_days, p.name
LIMIT {min(limit, 80)}
""".strip()
def _scenic_hint_terms(question: str) -> list[str]:
hints: list[str] = []
for scenic_name, aliases in SCENIC_HINT_ALIASES.items():
if scenic_name in question or any(alias in question for alias in aliases):
hints.append(scenic_name)
hints.extend(aliases)
out: list[str] = []
for hint in hints:
if hint and hint not in out:
out.append(hint)
return out
def _route_feature_terms(question: str) -> list[str]:
candidates = ("轻奢", "纯玩", "2+1", "头等舱", "精品", "小团", "私家团", "跟团", "亲子")
return [term for term in candidates if term in question]
def _route_multihop_cypher(question: str, limit: int) -> str | None:
terms = _route_name_terms(question)
if not terms:
return None
name_clauses: list[str] = []
for term in terms + _route_feature_terms(question):
quoted = _cypher_quote(term)
name_clauses.extend([
f"p.name CONTAINS '{quoted}'",
f"p.display_name CONTAINS '{quoted}'",
f"p.product_name CONTAINS '{quoted}'",
])
duration_clauses = [item for item in _route_duration_clauses(question) if item]
where_parts = [f"({' OR '.join(name_clauses)})"]
if duration_clauses:
where_parts.append(f"({' OR '.join(duration_clauses)})")
where_clause = " AND ".join(where_parts)
return f"""
MATCH (p:TourProduct)
WHERE {where_clause}
OPTIONAL MATCH (p)-[:PRODUCT_HAS_PRICE_PLAN]->(plan:ProductPricePlan)
OPTIONAL MATCH (p)-[:PRODUCT_HAS_DAY]->(d:ProductDay)-[:DAY_HAS_STOP]->(s:RouteStop)
OPTIONAL MATCH (s)-[:STOP_VISITS_ATTRACTION]->(a:ScenicAttraction)
OPTIONAL MATCH (a)-[near:ATTRACTION_NEARBY_HOTEL]-(h:Hotel)
RETURN
p.name AS product_name,
p.product_id AS product_id,
p.duration_days AS duration_days,
collect(DISTINCT {{
plan_name: plan.plan_name,
display_name: plan.display_name,
adult_price_min: plan.adult_price_min,
adult_price_max: plan.adult_price_max,
child_price_min: plan.child_price_min,
child_price_max: plan.child_price_max,
single_room_diff_min: plan.single_room_diff_min,
single_room_diff_max: plan.single_room_diff_max
}}) AS price_plans,
collect(DISTINCT {{
day: d.day_index,
sequence: s.sequence,
stop_name: s.name,
attraction_name: a.name
}}) AS route_points,
collect(DISTINCT {{
scenic_name: a.name,
stop_name: s.name,
hotel_name: h.name,
star_rating: h.star_rating,
city: h.city,
county: h.county,
address: h.address,
features: h.features,
driving_distance_km: near.driving_distance_km,
driving_time_minutes: near.driving_time_minutes,
recommend_rank: near.recommend_rank,
usage_note: near.usage_note,
remark: near.remark
}}) AS nearby_hotels
LIMIT {min(limit, 60)}
""".strip()
def _schema_snapshot(graph_name: str) -> dict[str, Any]:
now = time.monotonic()
cached = _SCHEMA_CACHE.get(graph_name)
if cached and now - cached[0] <= SCHEMA_CACHE_TTL_SECONDS:
return cached[1]
graph = _get_graph(graph_name)
labels: list[dict[str, Any]] = []
relations: list[dict[str, Any]] = []
samples: list[dict[str, Any]] = []
try:
for labels_value, count in graph.query(
"MATCH (n) RETURN labels(n) AS labels, count(n) AS count LIMIT 120"
).result_set:
label = labels_value[0] if labels_value else ""
labels.append({"label": label, "count": count})
except Exception:
labels = []
try:
for rel_type, count in graph.query(
"MATCH ()-[r]->() RETURN type(r) AS type, count(r) AS count LIMIT 160"
).result_set:
relations.append({"type": rel_type, "count": count})
except Exception:
relations = []
for item in labels[:12]:
label = item.get("label")
if not label or not re.match(r"^[A-Za-z_][A-Za-z0-9_]*$", str(label)):
continue
try:
rows = graph.query(f"MATCH (n:{label}) RETURN n LIMIT 1").result_set
except Exception:
continue
for row in rows:
node = row[0]
props = dict(getattr(node, "properties", None) or {})
samples.append({
"label": label,
"properties": {
str(k): _compact_scalar(v, 120)
for k, v in list(props.items())[:12]
},
})
snapshot = {"labels": labels[:80], "relations": relations[:100], "node_samples": samples[:16]}
_SCHEMA_CACHE[graph_name] = (now, snapshot)
return snapshot
def _list_texts(value: Any, limit: int, text_limit: int) -> list[str]:
if not isinstance(value, list):
return []
out: list[str] = []
for item in value[:limit]:
text = re.sub(r"\s+", " ", str(item or "")).strip()
if text:
out.append(text[:text_limit])
return out
def _is_route_list_question(question: str) -> bool:
return any(term in question for term in ("线路有哪些", "路线有哪些", "产品有哪些", "有哪些线路", "有哪些路线", "有哪些产品", "旅行车线路"))
def _route_items_from_rows(graph_result: dict[str, Any]) -> list[dict[str, Any]]:
items: list[dict[str, Any]] = []
seen: set[str] = set()
for row in graph_result.get("rows") or []:
if not isinstance(row, list) or not row:
continue
name = ""
duration = None
for value in row:
if isinstance(value, str) and value.strip():
name = value.strip()
break
for value in row:
if isinstance(value, (int, float)) and 0 < value <= 30:
duration = int(value)
break
if name and name not in seen:
seen.add(name)
items.append({"name": name, "duration_days": duration})
return items
def _money(value: Any) -> str:
if value in (None, ""):
return ""
try:
number = float(value)
except Exception:
return str(value)
if number.is_integer():
return str(int(number))
return f"{number:.2f}".rstrip("0").rstrip(".")
def _money_range(min_value: Any, max_value: Any, unit: str = "元/人") -> str:
low = _money(min_value)
high = _money(max_value)
if low and high and low != high:
return f"{low}-{high}{unit}"
if low or high:
return f"{low or high}{unit}"
return ""
def _price_plans_from_rows(graph_result: dict[str, Any]) -> list[dict[str, Any]]:
products: list[dict[str, Any]] = []
for row in graph_result.get("rows") or []:
if not isinstance(row, list) or len(row) < 4:
continue
product_name = str(row[0] or "").strip()
if not product_name:
continue
plans: list[dict[str, Any]] = []
for raw_plan in row[3] if isinstance(row[3], list) else []:
if not isinstance(raw_plan, dict):
continue
if not any(raw_plan.get(k) not in (None, "") for k in (
"plan_name", "display_name", "adult_price_min", "adult_price_max",
"child_price_min", "child_price_max", "single_room_diff_min",
"single_room_diff_max",
)):
continue
plans.append(raw_plan)
products.append({
"product_name": product_name,
"product_id": row[1] if len(row) > 1 else "",
"duration_days": row[2] if len(row) > 2 else None,
"plans": plans,
})
return products
def _route_list_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str] | None:
if not _is_route_list_question(question):
return None
items = _route_items_from_rows(graph_result)
if not items:
return None
lines = [
f"当前图谱共查到 {len(items)} 条旅行线路:",
*[
f"{idx}. {item['name']}" + (f"({item['duration_days']}日)" if item.get("duration_days") else "")
for idx, item in enumerate(items, start=1)
],
"如需进一步报价,请继续指定线路名称、出发日期、人数和住宿偏好。",
]
answer = "\n".join(lines)
preview_items = items[:10]
customer = "目前可查到 " + str(len(items)) + " 条旅行线路,主要包括:" + "、".join(
item["name"] + (f"({item['duration_days']}日)" if item.get("duration_days") else "")
for item in preview_items
)
if len(items) > len(preview_items):
customer += f"等,剩余 {len(items) - len(preview_items)} 条可继续筛选。"
else:
customer += "。"
customer += "您想看哪条线路的价格、行程或余位?"
return answer, customer
def _price_template_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str, list[dict[str, str]], list[dict[str, Any]]] | None:
products = _price_plans_from_rows(graph_result)
if not products:
return None
lines = [f"已在图谱中匹配到 {len(products)} 条相关线路报价:"]
evidence: list[dict[str, str]] = []
plans: list[dict[str, Any]] = []
for idx, product in enumerate(products[:8], start=1):
title = str(product["product_name"])
day_suffix = f"({product['duration_days']}日)" if product.get("duration_days") else ""
plan_lines: list[str] = []
for plan in product["plans"][:8]:
plan_name = str(plan.get("display_name") or plan.get("plan_name") or "默认报价档位")
adult = _money_range(plan.get("adult_price_min"), plan.get("adult_price_max"))
child = _money_range(plan.get("child_price_min"), plan.get("child_price_max"))
room = _money_range(plan.get("single_room_diff_min"), plan.get("single_room_diff_max"), "元")
parts = [plan_name]
if adult:
parts.append(f"成人 {adult}")
if child:
parts.append(f"儿童 {child}")
if room:
parts.append(f"单房差 {room}")
plan_lines.append(",".join(parts))
if not plan_lines:
plan_lines.append("图谱暂未录入明确报价档位,需要按团期核价")
lines.append(f"{idx}. {title}{day_suffix}:" + ";".join(plan_lines))
summary = ";".join(plan_lines)[:360]
evidence.append({
"type": "报价证据",
"name": title[:120],
"summary": summary,
"source": "FalkorDB 图查询 ProductPricePlan",
})
plans.append({
"rank": idx,
"label": "线路报价",
"plan_name": title,
"product_name": title,
"fit_score": max(60, 96 - idx * 3),
"match_reasons": ["命中线路名称/天数", "读取 ProductPricePlan 报价档位"],
"route_summary": f"{title}{day_suffix}",
"quote_summary": summary,
"variant_summary": summary,
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_price_template",
})
lines.append("以上为图谱报价档位,最终价格、余位、房型和儿童口径需要按具体出发日期二次核实。")
answer = "\n".join(lines)
first = products[0]
first_plans = first.get("plans") or []
if first_plans:
plan = first_plans[0]
adult = _money_range(plan.get("adult_price_min"), plan.get("adult_price_max"))
child = _money_range(plan.get("child_price_min"), plan.get("child_price_max"))
bits = [str(first["product_name"])]
if adult:
bits.append(f"成人参考 {adult}")
if child:
bits.append(f"儿童参考 {child}")
customer = ",".join(bits) + "。具体价格和余位需要按出发日期、人数、住宿档位再核实。"
else:
customer = f"已匹配到 {first['product_name']},但图谱暂未给出明确报价档位,需要按出发日期和人数核价。"
return answer, customer, evidence, plans
def _non_empty_maps(value: Any) -> list[dict[str, Any]]:
if not isinstance(value, list):
return []
out: list[dict[str, Any]] = []
for item in value:
if isinstance(item, dict) and any(v not in (None, "") for v in item.values()):
out.append(item)
return out
def _route_match_score(question: str, product: dict[str, Any]) -> int:
name = str(product.get("product_name") or "")
score = 0
for term in _route_name_terms(question):
if term and term in name:
score += 12
for term in _route_feature_terms(question):
if term and term in name:
score += 18
match = re.search(r"(\d+)\s*(?:日|天)", question)
if match and _safe_int(product.get("duration_days"), 0) == int(match.group(1)):
score += 15
for chinese, value in {"一": 1, "二": 2, "两": 2, "三": 3, "四": 4, "五": 5}.items():
if (f"{chinese}日" in question or f"{chinese}天" in question) and _safe_int(product.get("duration_days"), 0) == value:
score += 15
return score
def _route_multihop_products_from_rows(question: str, graph_result: dict[str, Any]) -> list[dict[str, Any]]:
products: list[dict[str, Any]] = []
for row in graph_result.get("rows") or []:
if not isinstance(row, list) or len(row) < 6:
continue
product = {
"product_name": str(row[0] or "").strip(),
"product_id": row[1] if len(row) > 1 else "",
"duration_days": row[2] if len(row) > 2 else None,
"price_plans": _non_empty_maps(row[3]),
"route_points": _non_empty_maps(row[4]),
"nearby_hotels": _non_empty_maps(row[5]),
}
if product["product_name"]:
product["match_score"] = _route_match_score(question, product)
products.append(product)
products.sort(key=lambda item: (item.get("match_score") or 0, -_safe_int(item.get("duration_days"), 99)), reverse=True)
return products
def _aggregate_price_text(plans: list[dict[str, Any]]) -> tuple[str, str, str]:
adult_lows: list[float] = []
adult_highs: list[float] = []
child_lows: list[float] = []
child_highs: list[float] = []
for plan in plans:
for bucket, key in ((adult_lows, "adult_price_min"), (adult_highs, "adult_price_max"), (child_lows, "child_price_min"), (child_highs, "child_price_max")):
try:
value = float(plan.get(key))
except Exception:
continue
if value > 0:
bucket.append(value)
adult = _money_range(min(adult_lows) if adult_lows else None, max(adult_highs) if adult_highs else None)
child = _money_range(min(child_lows) if child_lows else None, max(child_highs) if child_highs else None)
samples: list[str] = []
for plan in plans[:5]:
name = str(plan.get("display_name") or plan.get("plan_name") or "报价档位").strip()
adult_item = _money_range(plan.get("adult_price_min"), plan.get("adult_price_max"))
child_item = _money_range(plan.get("child_price_min"), plan.get("child_price_max"))
parts = [name]
if adult_item:
parts.append(f"成人 {adult_item}")
if child_item:
parts.append(f"儿童 {child_item}")
samples.append(",".join(parts))
return adult, child, ";".join(samples)
def _route_point_key(point: dict[str, Any]) -> tuple[int, int, str]:
return (
_safe_int(point.get("day"), 99),
_safe_int(point.get("sequence"), 99),
str(point.get("attraction_name") or point.get("stop_name") or ""),
)
def _is_non_scenic_route_node(name: str) -> bool:
compact = re.sub(r"\s+", "", name)
if any(term in compact for term in ("送机", "送站", "接机", "接站", "机场", "高铁站", "火车站", "散团", "集合")):
return True
if compact in {"贵阳", "安顺", "贵阳/安顺", "贵阳安顺"}:
return True
return bool(re.fullmatch(r"(贵阳|安顺|遵义|凯里|铜仁|荔波|西江)(/|、)(贵阳|安顺|遵义|凯里|铜仁|荔波|西江)", compact))
def _route_scenic_names(points: list[dict[str, Any]], limit: int = 12) -> list[str]:
out: list[str] = []
for point in sorted(points, key=_route_point_key):
name = str(point.get("attraction_name") or point.get("stop_name") or "").strip()
if _is_non_scenic_route_node(name):
continue
if name and name not in out:
out.append(name)
if len(out) >= limit:
break
return out
def _hotel_sort_key(item: dict[str, Any]) -> tuple[int, float, str]:
rank = _safe_int(item.get("recommend_rank"), 999)
try:
distance = float(item.get("driving_distance_km"))
except Exception:
distance = 9999.0
return rank, distance, str(item.get("hotel_name") or "")
def _filter_hotels(question: str, hotels: list[dict[str, Any]], limit: int = 8) -> list[dict[str, Any]]:
hints = _scenic_hint_terms(question)
out: list[dict[str, Any]] = []
seen: set[str] = set()
for item in sorted(hotels, key=_hotel_sort_key):
hotel_name = str(item.get("hotel_name") or "").strip()
if not hotel_name:
continue
scenic_text = " ".join(str(item.get(field) or "") for field in ("scenic_name", "stop_name", "county", "city", "address"))
if hints and not any(hint and hint in scenic_text for hint in hints):
continue
key = f"{item.get('scenic_name') or item.get('stop_name')}|{hotel_name}"
if key in seen:
continue
seen.add(key)
out.append(item)
if len(out) >= limit:
break
if out or hints:
return out
for item in sorted(hotels, key=_hotel_sort_key):
hotel_name = str(item.get("hotel_name") or "").strip()
if not hotel_name:
continue
key = f"{item.get('scenic_name') or item.get('stop_name')}|{hotel_name}"
if key in seen:
continue
seen.add(key)
out.append(item)
if len(out) >= limit:
break
return out
def _hotel_line(item: dict[str, Any]) -> str:
name = str(item.get("hotel_name") or "").strip()
scenic = str(item.get("scenic_name") or item.get("stop_name") or "相关景区").strip()
star = str(item.get("star_rating") or "").strip()
location = " ".join(str(item.get(field) or "").strip() for field in ("county", "address") if item.get(field)).strip()
distance = _money(item.get("driving_distance_km"))
minutes = _money(item.get("driving_time_minutes"))
parts = [name]
if star:
parts.append(star)
if distance:
parts.append(f"距{scenic}约{distance}km")
if minutes:
parts.append(f"车程约{minutes}分钟")
if location:
parts.append(location)
return ",".join(parts)
def _route_multihop_template_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str, list[dict[str, str]], list[dict[str, Any]]] | None:
products = _route_multihop_products_from_rows(question, graph_result)
if not products:
return None
product = products[0]
product_name = str(product["product_name"])
day_text = f"{product['duration_days']}天" if product.get("duration_days") else "天数待核实"
adult_price, child_price, price_samples = _aggregate_price_text(product["price_plans"])
scenic_names = _route_scenic_names(product["route_points"])
hotels = _filter_hotels(question, product["nearby_hotels"], limit=8)
price_parts = []
if adult_price:
price_parts.append(f"成人参考 {adult_price}")
if child_price:
price_parts.append(f"儿童参考 {child_price}")
price_text = ",".join(price_parts) or "图谱暂未录入明确报价档位"
scenic_text = "、".join(scenic_names) if scenic_names else "图谱暂未查到明确途经景区"
hotel_texts = [_hotel_line(item) for item in hotels[:6]]
if hotel_texts:
hotel_text = ";".join(hotel_texts)
elif _scenic_hint_terms(question):
hotel_text = "该景区附近酒店证据不足,需要按具体入住地再核实。"
else:
hotel_text = "请指定一个景区后,我可以继续筛选附近可住酒店。"
lines = [
f"已按百姓惠图谱匹配到线路:{product_name}。",
f"1. 游玩天数:{day_text}。",
f"2. 参考价格:{price_text}。",
f"3. 期间可去景区/节点:{scenic_text}。",
f"4. 附近酒店参考:{hotel_text}",
]
if price_samples:
lines.append(f"报价档位示例:{price_samples}。")
lines.append("以上为图谱证据口径,最终价格、房型、余位、入住酒店和景区政策需按出发日期二次核实。")
answer = "\n".join(lines)
customer_bits = [f"{product_name} 是 {day_text}线路"]
if price_text:
customer_bits.append(price_text)
if scenic_names:
customer_bits.append("可玩 " + "、".join(scenic_names[:5]))
if hotels:
customer_bits.append("附近可参考 " + "、".join(str(item.get("hotel_name") or "") for item in hotels[:3] if item.get("hotel_name")))
customer_reply = "亲," + ";".join(customer_bits) + "。具体价格、房型和余位需要按出发日期再确认。"
evidence = [
{
"type": "线路证据",
"name": product_name[:120],
"summary": f"{product_name};{day_text};{price_text}",
"source": "FalkorDB 图查询 TourProduct/ProductPricePlan",
},
{
"type": "行程景区证据",
"name": "途经景区",
"summary": scenic_text[:360],
"source": "FalkorDB 图查询 ProductDay/RouteStop/ScenicAttraction",
},
]
if hotel_texts:
evidence.append({
"type": "酒店证据",
"name": "景区附近酒店",
"summary": ";".join(hotel_texts[:4])[:360],
"source": "FalkorDB 图查询 ATTRACTION_NEARBY_HOTEL/Hotel",
})
plans = [{
"rank": 1,
"label": "线路多跳问答",
"plan_name": product_name,
"product_name": product_name,
"fit_score": 94,
"match_reasons": ["命中线路实体", "读取报价、天数、景区和附近酒店多跳证据"],
"route_summary": f"{product_name};{day_text};{scenic_text}",
"quote_summary": price_text,
"variant_summary": hotel_text,
"daily_itinerary": [
{
"day": point.get("day"),
"title": point.get("attraction_name") or point.get("stop_name"),
"summary": point.get("stop_name") or "",
}
for point in sorted(product["route_points"], key=_route_point_key)[:12]
if point.get("attraction_name") or point.get("stop_name")
],
"hotels": [
{
"name": item.get("hotel_name"),
"scenic_name": item.get("scenic_name") or item.get("stop_name"),
"summary": _hotel_line(item),
}
for item in hotels[:6]
],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_route_multihop_template",
}]
return answer, customer_reply, evidence, plans
def _valid_travel_item(item: Any) -> bool:
return isinstance(item, dict) and bool(str(item.get("name") or "").strip())
def _dedupe_travel_items(items: list[dict[str, Any]], limit: int = 18) -> list[dict[str, Any]]:
out: list[dict[str, Any]] = []
seen: set[str] = set()
for item in items:
if not _valid_travel_item(item):
continue
key = "|".join(str(item.get(field) or "") for field in ("name", "price", "adult_price", "unit", "status"))
if key in seen:
continue
seen.add(key)
out.append(item)
if len(out) >= limit:
break
return out
def _fee_products_from_rows(graph_result: dict[str, Any]) -> list[dict[str, Any]]:
products: list[dict[str, Any]] = []
for row in graph_result.get("rows") or []:
if not isinstance(row, list) or len(row) < 6:
continue
items: list[dict[str, Any]] = []
for group in row[3:6]:
if isinstance(group, list):
items.extend(item for item in group if isinstance(item, dict))
products.append({
"product_name": str(row[0] or "").strip(),
"product_id": row[1] if len(row) > 1 else "",
"duration_days": row[2] if len(row) > 2 else None,
"items": _dedupe_travel_items(items),
})
return [product for product in products if product["product_name"] and product["items"]]
def _item_price_text(item: dict[str, Any]) -> str:
price = _money(item.get("adult_price") or item.get("price"))
unit = str(item.get("unit") or "人").strip()
if price:
return f"{price}元/{unit}"
return "需核价"
def _fee_template_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str, list[dict[str, str]], list[dict[str, Any]]] | None:
products = _fee_products_from_rows(graph_result)
if not products:
return None
lines = [f"已在图谱中匹配到 {len(products)} 条线路的费用/小交通项目:"]
evidence: list[dict[str, str]] = []
plans: list[dict[str, Any]] = []
for idx, product in enumerate(products[:6], start=1):
title = product["product_name"]
day_suffix = f"({product['duration_days']}日)" if product.get("duration_days") else ""
item_texts: list[str] = []
for item in product["items"][:10]:
item_text = f"{item.get('name')}:{_item_price_text(item)}"
status = str(item.get("status") or item.get("default_status") or "").strip()
if status:
item_text += f"({status})"
item_texts.append(item_text)
summary = ";".join(item_texts)
lines.append(f"{idx}. {title}{day_suffix}:" + summary)
evidence.append({
"type": "费用证据",
"name": title[:120],
"summary": summary[:360],
"source": "FalkorDB 图查询 TravelItem",
})
plans.append({
"rank": idx,
"label": "费用/小交通",
"plan_name": title,
"product_name": title,
"fit_score": max(60, 94 - idx * 3),
"match_reasons": ["命中线路名称/天数", "读取 TravelItem 费用项目"],
"route_summary": f"{title}{day_suffix}",
"quote_summary": summary[:360],
"variant_summary": summary[:360],
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_fee_template",
})
lines.append("以上项目按图谱记录展示,景区门票、小交通和可选项目需按团期及景区政策二次核实。")
answer = "\n".join(lines)
first = products[0]
preview = ";".join(f"{item.get('name')} {_item_price_text(item)}" for item in first["items"][:4])
customer = f"{first['product_name']} 相关费用项目包括:{preview}。具体是否必含、是否可选和儿童政策需要按团期再核实。"
return answer, customer, evidence, plans
def _vehicle_items_from_rows(graph_result: dict[str, Any]) -> list[dict[str, Any]]:
items: list[dict[str, Any]] = []
for row in graph_result.get("rows") or []:
if not isinstance(row, list) or not row:
continue
name = str(row[0] or "").strip()
if not name:
continue
items.append({
"name": name,
"type": row[1] if len(row) > 1 else "",
"subtype": row[2] if len(row) > 2 else "",
"price": row[3] if len(row) > 3 else None,
"unit": row[4] if len(row) > 4 else "",
"raw_evidence": row[5] if len(row) > 5 else "",
})
return _dedupe_travel_items(items, limit=30)
def _vehicle_template_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str, list[dict[str, str]], list[dict[str, Any]]] | None:
items = _vehicle_items_from_rows(graph_result)
if not items:
return None
lines = [
f"当前图谱查到 {len(items)} 个用车/车型相关项目:",
*[
f"{idx}. {item['name']}" + (f":{_item_price_text(item)}" if item.get("price") not in (None, "") else "")
for idx, item in enumerate(items, start=1)
],
"具体车型是否可用、座位数、车辆档位和结算口径需要按团期及供应商确认。",
]
answer = "\n".join(lines)
preview = "、".join(
item["name"] + (f"({_item_price_text(item)})" if item.get("price") not in (None, "") else "")
for item in items[:8]
)
customer = f"目前图谱里可查到 {len(items)} 个用车相关项目,主要包括:{preview}。具体车辆安排需按出发日期、人数和供应商确认。"
evidence = [
{
"type": "用车证据",
"name": item["name"][:120],
"summary": f"{item.get('subtype') or item.get('type') or '用车项目'};参考价:{_item_price_text(item)}",
"source": "FalkorDB 图查询 TravelItem",
}
for item in items[:12]
]
plans = [
{
"rank": idx,
"label": "用车项目",
"plan_name": item["name"],
"product_name": item["name"],
"fit_score": max(60, 92 - idx),
"match_reasons": ["用车/车型快查", "读取 TravelItem"],
"route_summary": item["name"],
"quote_summary": _item_price_text(item),
"variant_summary": str(item.get("subtype") or ""),
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_vehicle_template",
}
for idx, item in enumerate(items[:4], start=1)
]
return answer, customer, evidence, plans
def _context_products_from_rows(graph_result: dict[str, Any]) -> list[dict[str, Any]]:
products: list[dict[str, Any]] = []
for row in graph_result.get("rows") or []:
if not isinstance(row, list) or len(row) < 5:
continue
name = str(row[0] or "").strip()
if not name:
continue
price_plans = [item for item in row[3] if isinstance(item, dict)] if isinstance(row[3], list) else []
route_points = [item for item in row[4] if isinstance(item, dict)] if isinstance(row[4], list) else []
min_price = None
for plan in price_plans:
value = plan.get("adult_price_min")
if isinstance(value, (int, float)):
min_price = value if min_price is None else min(min_price, value)
scenic_names = _route_scenic_names(route_points, limit=30)
attraction_count = len(scenic_names)
products.append({
"product_name": name,
"product_id": row[1] if len(row) > 1 else "",
"duration_days": row[2] if len(row) > 2 else None,
"adult_price_min": min_price,
"price_plans": price_plans[:8],
"route_points": route_points[:16],
"scenic_names": scenic_names,
"attraction_count": attraction_count,
})
return products
def _condition_route_score(question: str, product: dict[str, Any]) -> tuple[int, int, int]:
name = str(product.get("product_name") or "")
duration = _safe_int(product.get("duration_days"), 99)
price = product.get("adult_price_min")
price_int = _safe_int(price, 999999) if price not in (None, "") else 999999
scenic_count = _safe_int(product.get("attraction_count"), len(product.get("scenic_names") or []))
scenic_text = " ".join(product.get("scenic_names") or [])
slots = _condition_slots(question)
score = 50
if slots["family"]:
if any(term in name for term in ("亲子", "轻奢", "纯玩", "小包团", "2+1", "私家")):
score += 18
if duration <= 3:
score += 10
elif duration <= 4:
score += 6
if slots["comfort"]:
if any(term in name for term in ("轻奢", "纯玩", "2+1", "私家", "头等舱")):
score += 20
if duration <= 3:
score += 12
if 0 < scenic_count <= 5:
score += 6
if duration >= 5:
score -= 8
if "梵净山" in scenic_text:
score -= 8
if slots["budget"]:
if price_int <= 800:
score += 18
elif price_int <= 1200:
score += 10
if any(term in name for term in ("购物", "特价购物")):
score -= 12
return score, -duration, -price_int
def _condition_route_template_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str, list[dict[str, str]], list[dict[str, Any]]] | None:
products = _context_products_from_rows(graph_result)
if not products:
return None
products.sort(key=lambda item: _condition_route_score(question, item), reverse=True)
selected = products[:5]
lines = ["按“老人/小孩/轻松/预算”等条件从图谱中筛出以下候选线路:"]
evidence: list[dict[str, str]] = []
plans: list[dict[str, Any]] = []
for idx, product in enumerate(selected, start=1):
name = product["product_name"]
duration = f"{product['duration_days']}日" if product.get("duration_days") else "天数待核实"
price = _money(product.get("adult_price_min"))
price_text = f"成人起价约 {price}元/人" if price else "价格需按团期核实"
scenic_names = product.get("scenic_names") or _route_scenic_names(product.get("route_points") or [], limit=8)
scenic_text = "、".join(scenic_names[:5]) if scenic_names else "景区信息待核实"
reasons: list[str] = []
if any(term in name for term in ("轻奢", "纯玩", "2+1", "私家", "小包团")):
reasons.append("线路名称体现舒适/纯玩/小团倾向")
if _safe_int(product.get("duration_days"), 99) <= 3:
reasons.append("天数较短,行程压力相对低")
if price:
reasons.append("图谱有报价起价可参考")
reason_text = ";".join(reasons[:3]) or "图谱有线路、价格和行程证据"
line = f"{idx}. {name}({duration}):{price_text};可玩 {scenic_text};推荐依据:{reason_text}"
lines.append(line)
evidence.append({
"type": "条件筛选证据",
"name": name[:120],
"summary": line[:360],
"source": "FalkorDB 图查询 TourProduct/ProductPricePlan/ScenicAttraction",
})
plans.append({
"rank": idx,
"label": "条件推荐线路",
"plan_name": name,
"product_name": name,
"fit_score": max(60, min(96, _condition_route_score(question, product)[0])),
"match_reasons": reasons or ["条件筛选命中", "读取线路报价和景区证据"],
"route_summary": f"{name}({duration}):{scenic_text}",
"quote_summary": price_text,
"variant_summary": reason_text,
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_condition_route_advice",
})
best = selected[0]
best_price = _money(best.get("adult_price_min"))
best_price_text = f",成人起价约 {best_price}元/人" if best_price else ""
best_duration = f"{best['duration_days']}日" if best.get("duration_days") else "天数待核实"
customer = f"亲,带老人小孩又希望轻松,可以优先看 {best['product_name']}({best_duration}{best_price_text})。它在图谱里有线路和景区证据,最终还要按出发日期、人数和住宿档位核实。"
lines.append("这些推荐是基于图谱字段做的初筛,不等于最终承诺;老人/儿童还需要二次确认行程强度、用车、住宿和余位。")
return "\n".join(lines), customer, evidence, plans
def _route_compare_suitability_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str, list[dict[str, str]], list[dict[str, Any]]] | None:
products = _context_products_from_rows(graph_result)
if not products:
return None
compare_terms = _route_name_terms(question)
grouped: list[dict[str, Any]] = []
for term in compare_terms:
candidates = [item for item in products if term in str(item.get("product_name") or "")]
if not candidates:
continue
other_terms = [other for other in compare_terms if other != term]
pure_candidates = [
item for item in candidates
if not any(other in str(item.get("product_name") or "") for other in other_terms)
]
pool = pure_candidates or candidates
pool.sort(key=lambda item: _condition_route_score(question, item), reverse=True)
best = pool[0]
grouped.append({
"term": term,
"best": best,
"score": _condition_route_score(question, best)[0],
})
if len(grouped) < 2:
return None
grouped.sort(key=lambda item: item["score"], reverse=True)
winner = grouped[0]
lines = ["按老人/儿童/轻松度条件,用图谱中的天数、价格和景区节点做适配度对比:"]
evidence: list[dict[str, str]] = []
plans: list[dict[str, Any]] = []
for idx, group in enumerate(grouped, start=1):
product = group["best"]
name = product["product_name"]
duration = f"{product['duration_days']}日" if product.get("duration_days") else "天数待核实"
price = _money(product.get("adult_price_min"))
price_text = f"成人起价约 {price}元/人" if price else "价格需按团期核实"
scenic_names = product.get("scenic_names") or []
scenic_text = "、".join(scenic_names[:6]) if scenic_names else "景区信息待核实"
reasons: list[str] = []
if _safe_int(product.get("duration_days"), 99) <= 3:
reasons.append("天数短,行程压力相对低")
if any(term in name for term in ("轻奢", "2+1", "纯玩", "头等舱", "私家")):
reasons.append("名称体现舒适/纯玩/车型优势")
if "梵净山" in scenic_text:
reasons.append("包含梵净山,老人小孩需评估体力")
if price:
reasons.append("有报价起价可参考")
reason_text = ";".join(reasons[:4]) or "有线路、价格和景区证据"
line = f"{idx}. {group['term']}:代表线路 {name}({duration}),{price_text};景区:{scenic_text};判断:{reason_text}"
lines.append(line)
evidence.append({
"type": "适配度对比证据",
"name": str(group["term"])[:120],
"summary": line[:360],
"source": "FalkorDB 图查询 TourProduct/ProductPricePlan/ScenicAttraction",
})
plans.append({
"rank": idx,
"label": "线路适配度对比",
"plan_name": str(group["term"]),
"product_name": name,
"fit_score": max(60, min(96, int(group["score"]))),
"match_reasons": reasons or ["条件适配评分", "读取线路报价和景区证据"],
"route_summary": line,
"quote_summary": price_text,
"variant_summary": reason_text,
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_route_compare_suitability",
})
winner_product = winner["best"]
customer = f"亲,按老人小孩和不要太累的条件,优先建议 {winner['term']},代表线路可看 {winner_product['product_name']}。它在图谱评分里更偏轻松/舒适;最终还要按出发日期、人数、住宿和老人儿童体力二次确认。"
lines.append("该结论是图谱条件评分,不是最终承诺;涉及老人儿童时建议再核实步行强度、用车、住宿和景区政策。")
return "\n".join(lines), customer, evidence, plans
def _route_compare_scenic_count_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str, list[dict[str, str]], list[dict[str, Any]]] | None:
products = _context_products_from_rows(graph_result)
if not products:
return None
compare_terms = _route_name_terms(question)
grouped: list[dict[str, Any]] = []
for term in compare_terms:
candidates = [item for item in products if term in str(item.get("product_name") or "")]
if not candidates:
continue
other_terms = [other for other in compare_terms if other != term]
pure_candidates = [
item for item in candidates
if not any(other in str(item.get("product_name") or "") for other in other_terms)
]
pool = pure_candidates or candidates
pool.sort(key=lambda item: (_safe_int(item.get("attraction_count"), 0), _route_match_score(question, item)), reverse=True)
best = pool[0]
grouped.append({
"term": term,
"best": best,
"count": _safe_int(best.get("attraction_count"), 0),
"candidates": pool[:4],
})
if grouped:
selected = [item["best"] for item in grouped]
else:
products.sort(key=lambda item: (_route_match_score(question, item), _safe_int(item.get("attraction_count"), 0)), reverse=True)
selected = products[:8]
max_count = max(_safe_int(item.get("attraction_count"), 0) for item in selected)
winners = [item for item in selected if _safe_int(item.get("attraction_count"), 0) == max_count and max_count > 0]
lines = ["按当前图谱的 ScenicAttraction 证据统计:"]
evidence: list[dict[str, str]] = []
plans: list[dict[str, Any]] = []
if grouped:
for idx, group in enumerate(grouped, start=1):
product = group["best"]
scenic_names = product.get("scenic_names") or _route_scenic_names(product.get("route_points") or [], limit=10)
count = _safe_int(product.get("attraction_count"), len(scenic_names))
scenic_text = "、".join(scenic_names[:8]) if scenic_names else "暂无明确景区证据"
line = f"{idx}. {group['term']}:代表线路 {product['product_name']},图谱命中 {count} 个景区/游玩节点,包含 {scenic_text}"
lines.append(line)
evidence.append({
"type": "景区数量对比证据",
"name": str(group["term"])[:120],
"summary": line[:360],
"source": "FalkorDB 图查询 ProductDay/RouteStop/ScenicAttraction",
})
plans.append({
"rank": idx,
"label": "线路景点数量对比",
"plan_name": str(group["term"]),
"product_name": product["product_name"],
"fit_score": max(60, 92 - idx * 3),
"match_reasons": ["按用户提到的线路词分组", "统计 ScenicAttraction 去重数量"],
"route_summary": line,
"quote_summary": "",
"variant_summary": scenic_text,
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_route_compare_scenic_count",
})
else:
for idx, product in enumerate(selected, start=1):
scenic_names = product.get("scenic_names") or _route_scenic_names(product.get("route_points") or [], limit=10)
count = _safe_int(product.get("attraction_count"), len(scenic_names))
scenic_text = "、".join(scenic_names[:8]) if scenic_names else "暂无明确景区证据"
line = f"{idx}. {product['product_name']}:图谱命中 {count} 个景区/游玩节点,包含 {scenic_text}"
lines.append(line)
evidence.append({
"type": "景区数量对比证据",
"name": product["product_name"][:120],
"summary": line[:360],
"source": "FalkorDB 图查询 ProductDay/RouteStop/ScenicAttraction",
})
plans.append({
"rank": idx,
"label": "线路景点数量对比",
"plan_name": product["product_name"],
"product_name": product["product_name"],
"fit_score": max(60, 92 - idx * 3),
"match_reasons": ["线路名称命中", "统计 ScenicAttraction 去重数量"],
"route_summary": line,
"quote_summary": "",
"variant_summary": scenic_text,
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_route_compare_scenic_count",
})
if winners:
if len(winners) > 1 and grouped:
winner_names = "、".join(
group["term"] for group in grouped if group["best"] in winners
)
customer = f"从当前图谱景区数量看,{winner_names} 最高都约 {max_count} 个景区/游玩节点,数量接近;建议再结合天数、价格和行程强度选择。"
elif grouped:
winner_name = next((group["term"] for group in grouped if group["best"] in winners), "")
customer = f"从当前图谱景区数量看,{winner_name or winners[0]['product_name']}的代表线路命中的景区/游玩节点更多(约 {max_count} 个)。具体还要结合出发日期、行程强度和实际游玩安排确认。"
else:
winner_names = "、".join(item["product_name"] for item in winners[:2])
customer = f"从当前图谱景区数量看,{winner_names} 命中的景区/游玩节点更多(约 {max_count} 个)。具体还要结合出发日期、行程强度和实际游玩安排确认。"
else:
customer = "当前图谱没有足够的景区数量证据做可靠比较,建议补充具体线路名称或让后台核实行程明细。"
lines.append("注意:这里比较的是图谱已结构化的景区/游玩节点数量,不代表实际游玩时长或体验强弱。")
return "\n".join(lines), customer, evidence, plans
def _context_template_fallback_answer(question: str, graph_result: dict[str, Any]) -> tuple[str, str, list[dict[str, str]], list[dict[str, Any]]]:
products = _context_products_from_rows(graph_result)
products.sort(key=lambda item: (item.get("adult_price_min") is None, item.get("adult_price_min") or 999999, item.get("duration_days") or 99))
evidence: list[dict[str, str]] = []
plans: list[dict[str, Any]] = []
lines = ["已按线路实体查到以下可对比方案:"]
for idx, product in enumerate(products[:6], start=1):
price = _money(product.get("adult_price_min"))
day_suffix = f"({product['duration_days']}日)" if product.get("duration_days") else ""
price_text = f",成人参考起价 {price}元/人" if price else ",价格需按团期核实"
line = f"{idx}. {product['product_name']}{day_suffix}{price_text}"
lines.append(line)
evidence.append({
"type": "线路对比证据",
"name": product["product_name"][:120],
"summary": line[:360],
"source": "FalkorDB 图查询 TourProduct/ProductPricePlan",
})
plans.append({
"rank": idx,
"label": "线路对比",
"plan_name": product["product_name"],
"product_name": product["product_name"],
"fit_score": max(60, 90 - idx * 3),
"match_reasons": ["命中线路实体", "读取报价和行程上下文"],
"route_summary": line,
"quote_summary": line,
"variant_summary": "",
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "deterministic_context_fallback",
})
if products:
best = products[0]
customer = f"从图谱报价看,预算优先可以先看 {best['product_name']}。如果有老人和小孩,还需要结合行程强度、住宿档位和出发日期确认,价格和余位需二次核实。"
else:
customer = "当前已进入线路对比查询,但图谱证据不足,建议补充具体线路、出发日期、人数和预算后再推荐。"
lines.append("建议最终按出发日期、人数、住宿档位、老人/儿童体力情况和供应商余位二次核实。")
return "\n".join(lines), customer, evidence, plans
async def _deterministic_context_response(
question: str,
graph_name: str,
*,
limit: int,
customer_context: dict[str, Any] | None,
started_at: float,
intent: dict[str, Any],
intent_ms: int,
) -> dict[str, Any] | None:
if intent.get("route") != "deterministic_context_llm_answer":
return None
if intent.get("intent") == "condition_route_advice":
cypher = _condition_context_cypher(limit)
else:
cypher = _route_context_cypher(question, limit) or ""
if not cypher:
return None
graph_started = time.perf_counter()
try:
graph_result = await asyncio.to_thread(_run_cypher, graph_name, cypher, limit)
except Exception:
return None
graph_ms = max(1, round((time.perf_counter() - graph_started) * 1000))
if not graph_result.get("row_count"):
return None
planner = intent.get("_planner_trace") if isinstance(intent.get("_planner_trace"), dict) else {}
planner_ms = _safe_int(planner.get("latency_ms"), 0)
answer_ms = 0
llm_error = ""
answer_data: dict[str, Any] = {}
answer_client = await _graph_qa_client(max_tokens=600)
if answer_client is not None:
payload = {
"question": question,
"graph_name": graph_name,
"intent": intent,
"answer_focus": "基于线路报价、天数、景区数量和行程证据做 3 句话以内建议;证据不足必须说明需二次核实。",
"cypher": cypher,
"graph_result": {
"columns": graph_result["columns"],
"row_count": graph_result["row_count"],
"rows": graph_result["rows"][:6],
"nodes": [],
"relationships": [],
},
"customer_context": customer_context or {},
}
answer_started = time.perf_counter()
try:
answer_data = await _chat_json_timed(
answer_client,
GRAPH_QA_ANSWER_SYS,
json.dumps(payload, ensure_ascii=False),
timeout_seconds=GRAPH_QA_ANSWER_TIMEOUT_SECONDS,
attempts=1,
)
except asyncio.TimeoutError:
llm_error = "answer_timeout"
except Exception as exc: # noqa: BLE001
llm_error = str(exc)[:260]
answer_ms = max(1, round((time.perf_counter() - answer_started) * 1000))
if answer_data:
answer = str(answer_data.get("answer") or answer_data.get("customer_reply") or "").strip()
customer_reply = str(answer_data.get("customer_reply") or answer).strip()
evidence = _evidence_cards(answer_data, graph_result)
plans = _plans_from_evidence(evidence, answer)
confidence = _safe_confidence(answer_data.get("confidence"), 0.72)
followups = _list_texts(answer_data.get("follow_up_questions"), 4, 120)
risk_notes = _list_texts(answer_data.get("risk_notes"), 4, 160)
else:
answer, customer_reply, evidence, plans = _context_template_fallback_answer(question, graph_result)
confidence = 0.68
followups = ["请提供出发日期和人数。", "老人和小孩对行程强度有什么要求?", "预算范围和住宿档位是多少?"]
risk_notes = ["推荐仅基于图谱已有报价和线路信息,最终需按团期、余位和供应商政策二次核实。"]
latency_ms = max(1, round((time.perf_counter() - started_at) * 1000))
return {
"question": question,
"graph_name": graph_name,
"answer": answer,
"customer_reply": customer_reply,
"copy_text": customer_reply or answer,
"plans": plans,
"evidence": evidence,
"sales_scripts": [],
"follow_up_questions": followups,
"risk_notes": risk_notes,
"confidence": confidence,
"graph_result": graph_result,
"trace": {
"method": "deterministic_context_llm_answer_v1",
"query_source": "deterministic_context_cypher",
"rule_query_used": False,
"llm_used": bool(answer_data),
"llm_error": llm_error,
"cache_hit": False,
"cache": {"response": False, "cypher": False},
"generated_cypher": "",
"effective_cypher": cypher,
"cypher_cache_hit": False,
"fallback_cypher": "",
"fallback_query_used": False,
"cypher_reason": intent.get("reason") or "",
"intent": intent.get("intent") or "",
"llm_planner_used": bool(planner.get("used")),
"llm_planner": planner,
"graph_qa_intent": intent,
"routing_strategy": "intent_context_query_then_llm_answer",
"cypher_repaired": False,
"first_query_error": "",
"row_count": graph_result["row_count"],
"node_count": len(graph_result["nodes"]),
"relationship_count": len(graph_result["relationships"]),
"latency_ms": latency_ms,
"stage_timings_ms": {
"intent_classification": intent_ms,
"llm_planner": planner_ms,
"schema": 0,
"cypher_generation": 0,
"graph_query": graph_ms,
"fallback_graph_query": 0,
"answer_synthesis": answer_ms,
},
"performance_target_ms": 1200,
"response_mode": "deterministic_context_llm_answer",
"response_mode_label": "固定取证 + LLM 证据回答",
"graph_capabilities_used": ["意图识别", "固定 Cypher 取证", "LLM 证据回答"],
"retrieval_summary": {
"cypher": cypher,
"llm_generated_cypher": "",
"fallback_query_used": False,
"rows": graph_result["row_count"],
"nodes": len(graph_result["nodes"]),
"relationships": len(graph_result["relationships"]),
},
},
}
def _cache_key(question: str, graph_name: str, limit: int, customer_context: dict[str, Any] | None) -> str:
context_key = json.dumps(customer_context or {}, ensure_ascii=False, sort_keys=True, default=str)[:800]
normalized_question = re.sub(r"\s+", "", question.strip())
return f"{graph_name}|{limit}|{normalized_question}|{context_key}"
def _get_cached_response(cache_key: str, started_at: float) -> dict[str, Any] | None:
cached = _QA_RESPONSE_CACHE.get(cache_key)
if not cached:
return None
cached_at, payload = cached
if time.monotonic() - cached_at > QA_RESPONSE_CACHE_TTL_SECONDS:
_QA_RESPONSE_CACHE.pop(cache_key, None)
return None
response = copy.deepcopy(payload)
latency_ms = max(1, round((time.perf_counter() - started_at) * 1000))
trace = response.setdefault("trace", {})
trace["cache_hit"] = True
cache = trace.setdefault("cache", {})
cache["response"] = True
trace["latency_ms"] = latency_ms
trace["stage_timings_ms"] = {
**(trace.get("stage_timings_ms") or {}),
"cache_lookup": latency_ms,
}
return response
def _set_cached_response(cache_key: str, response: dict[str, Any]) -> None:
if len(_QA_RESPONSE_CACHE) > 256:
oldest_key = min(_QA_RESPONSE_CACHE, key=lambda key: _QA_RESPONSE_CACHE[key][0])
_QA_RESPONSE_CACHE.pop(oldest_key, None)
_QA_RESPONSE_CACHE[cache_key] = (time.monotonic(), copy.deepcopy(response))
def _cypher_cache_key(question: str, graph_name: str, limit: int, intent: dict[str, Any]) -> str:
normalized = _normalized_question(question)
return f"{graph_name}|{limit}|{intent.get('intent') or 'general'}|{normalized}"
def _get_cached_cypher(cache_key: str) -> dict[str, Any] | None:
cached = _CYPHER_CACHE.get(cache_key)
if not cached:
return None
cached_at, payload = cached
if time.monotonic() - cached_at > CYPHER_CACHE_TTL_SECONDS:
_CYPHER_CACHE.pop(cache_key, None)
return None
return copy.deepcopy(payload)
def _set_cached_cypher(cache_key: str, decision: dict[str, Any]) -> None:
if len(_CYPHER_CACHE) > 512:
oldest_key = min(_CYPHER_CACHE, key=lambda key: _CYPHER_CACHE[key][0])
_CYPHER_CACHE.pop(oldest_key, None)
_CYPHER_CACHE[cache_key] = (time.monotonic(), copy.deepcopy(decision))
async def _fast_graph_response(
question: str,
graph_name: str,
*,
limit: int,
customer_context: dict[str, Any] | None,
started_at: float,
intent: dict[str, Any],
intent_ms: int,
) -> dict[str, Any] | None:
if intent.get("route") != "template_graph_query":
return None
fast_kind = str(intent.get("intent") or "")
cypher = ""
if fast_kind == "route_catalog":
cypher = _route_catalog_cypher(limit)
elif fast_kind == "condition_route_advice":
cypher = _condition_context_cypher(limit)
elif fast_kind == "route_compare_suitability":
cypher = _route_context_cypher(question, limit) or ""
elif fast_kind == "route_compare_scenic_count":
cypher = _route_context_cypher(question, limit) or ""
elif fast_kind == "route_multihop_detail":
cypher = _route_multihop_cypher(question, limit) or ""
elif fast_kind == "price_quote":
cypher = _price_fallback_cypher(question, limit) or ""
elif fast_kind == "fee_detail":
cypher = _fee_fallback_cypher(question, limit) or ""
elif fast_kind == "vehicle_catalog":
cypher = _vehicle_catalog_cypher(limit)
if not cypher:
return None
graph_started = time.perf_counter()
try:
graph_result = await asyncio.to_thread(_run_cypher, graph_name, cypher, limit)
except Exception:
return None
graph_ms = max(1, round((time.perf_counter() - graph_started) * 1000))
if not graph_result.get("row_count"):
return None
planner = intent.get("_planner_trace") if isinstance(intent.get("_planner_trace"), dict) else {}
planner_ms = _safe_int(planner.get("latency_ms"), 0)
if fast_kind == "route_catalog":
route_answer = _route_list_answer(question, graph_result)
if not route_answer:
return None
answer, customer_reply = route_answer
items = _route_items_from_rows(graph_result)
evidence = [
{
"type": "线路清单",
"name": item["name"],
"summary": f"{item['name']}" + (f"({item['duration_days']}日)" if item.get("duration_days") else ""),
"source": "FalkorDB 图查询 TourProduct",
}
for item in items[:12]
]
plans = [
{
"rank": idx,
"label": "线路清单",
"plan_name": item["name"],
"product_name": item["name"],
"fit_score": max(60, 95 - idx),
"match_reasons": ["线路清单快查", "读取 TourProduct"],
"route_summary": item["name"],
"quote_summary": "可继续按线路名称、日期和人数查询报价。",
"variant_summary": "",
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "fast_route_catalog",
}
for idx, item in enumerate(items[:4], start=1)
]
followups = ["要查哪条线路的价格?", "请提供出发日期、人数和住宿偏好。", "是否需要按天数或景点筛选线路?"]
confidence = 0.9
elif fast_kind == "condition_route_advice":
condition_answer = _condition_route_template_answer(question, graph_result)
if not condition_answer:
return None
answer, customer_reply, evidence, plans = condition_answer
followups = ["请提供出发日期和人数。", "老人或小孩是否有少走路要求?", "预算和住宿档位大概是多少?"]
confidence = 0.78
elif fast_kind == "route_compare_suitability":
suitability_answer = _route_compare_suitability_answer(question, graph_result)
if not suitability_answer:
return None
answer, customer_reply, evidence, plans = suitability_answer
followups = ["请提供出发日期和人数。", "老人小孩是否能接受爬山或较多步行?", "是否需要同时比较价格?"]
confidence = 0.8
elif fast_kind == "route_compare_scenic_count":
compare_answer = _route_compare_scenic_count_answer(question, graph_result)
if not compare_answer:
return None
answer, customer_reply, evidence, plans = compare_answer
followups = ["是否还要比较价格和行程强度?", "请提供出发日期和人数。", "是否需要按老人/儿童适配度再筛选?"]
confidence = 0.82
elif fast_kind == "route_multihop_detail":
multihop_answer = _route_multihop_template_answer(question, graph_result)
if not multihop_answer:
return None
answer, customer_reply, evidence, plans = multihop_answer
followups = ["请提供出发日期和人数。", "是否需要指定住宿档位或房型?", "要继续核实哪个景区附近酒店?"]
confidence = 0.88
elif fast_kind == "price_quote":
price_answer = _price_template_answer(question, graph_result)
if not price_answer:
return None
answer, customer_reply, evidence, plans = price_answer
followups = ["请提供出发日期和人数。", "是否有住宿档位要求?", "是否需要同步查询儿童价和单房差?"]
confidence = 0.86
elif fast_kind == "fee_detail":
fee_answer = _fee_template_answer(question, graph_result)
if not fee_answer:
return None
answer, customer_reply, evidence, plans = fee_answer
followups = ["要按哪一天或哪个景区核实费用?", "是否需要区分必含和自费项目?", "请提供出发日期方便核实景区政策。"]
confidence = 0.84
elif fast_kind == "vehicle_catalog":
vehicle_answer = _vehicle_template_answer(question, graph_result)
if not vehicle_answer:
return None
answer, customer_reply, evidence, plans = vehicle_answer
followups = ["请提供出发日期和人数。", "是否需要指定车型或座位布局?", "是否需要核算包车费用?"]
confidence = 0.82
else:
return None
latency_ms = max(1, round((time.perf_counter() - started_at) * 1000))
return {
"question": question,
"graph_name": graph_name,
"answer": answer,
"customer_reply": customer_reply,
"copy_text": customer_reply or answer,
"plans": plans,
"evidence": evidence,
"sales_scripts": [],
"follow_up_questions": followups,
"risk_notes": ["价格、余位、房型、车辆和景区政策以具体团期及供应商二次核实为准。"],
"confidence": confidence,
"graph_result": graph_result,
"trace": {
"method": f"fast_{fast_kind}_graph_template_v1",
"query_source": "deterministic_cypher_template",
"rule_query_used": False,
"llm_used": False,
"cache_hit": False,
"cache": {"response": False, "cypher": False},
"generated_cypher": "",
"effective_cypher": cypher,
"fallback_cypher": "",
"fallback_query_used": False,
"cypher_reason": "高频客服问题命中固定只读图查询模板,跳过 LLM 生成 Cypher 和 LLM 答案组织。",
"intent": fast_kind,
"llm_planner_used": bool(planner.get("used")),
"llm_planner": planner,
"graph_qa_intent": intent,
"routing_strategy": "intent_template_first",
"cypher_repaired": False,
"first_query_error": "",
"row_count": graph_result["row_count"],
"node_count": len(graph_result["nodes"]),
"relationship_count": len(graph_result["relationships"]),
"latency_ms": latency_ms,
"stage_timings_ms": {
"intent_classification": intent_ms,
"llm_planner": planner_ms,
"schema": 0,
"cypher_generation": 0,
"graph_query": graph_ms,
"fallback_graph_query": 0,
"answer_synthesis": 0,
},
"performance_target_ms": 1200,
"response_mode": "fast_graph_template",
"response_mode_label": "高频客服问题快查",
"graph_capabilities_used": ["固定 Cypher 模板", "FalkorDB 只读查询", "客服话术模板"],
"retrieval_summary": {
"cypher": cypher,
"llm_generated_cypher": "",
"fallback_query_used": False,
"rows": graph_result["row_count"],
"nodes": len(graph_result["nodes"]),
"relationships": len(graph_result["relationships"]),
},
},
}
def _evidence_cards(answer_data: dict[str, Any], graph_result: dict[str, Any]) -> list[dict[str, str]]:
cards: list[dict[str, str]] = []
raw_cards = answer_data.get("evidence_cards")
if isinstance(raw_cards, list):
for item in raw_cards[:6]:
if not isinstance(item, dict):
continue
title = str(item.get("title") or item.get("name") or "").strip()
summary = str(item.get("summary") or item.get("detail") or "").strip()
source = str(item.get("source") or "FalkorDB 图查询").strip()
if title or summary:
cards.append({"type": "图查询证据", "name": title[:120], "summary": summary[:360], "source": source[:120]})
for node in (graph_result.get("nodes") or [])[:6]:
if len(cards) >= 8:
break
props = node.get("properties") or {}
summary = ";".join(
f"{k}: {v}" for k, v in list(props.items())[:6] if v not in (None, "")
)
cards.append({
"type": "图谱节点",
"name": str(node.get("title") or "")[:120],
"summary": summary[:360],
"source": "FalkorDB 图查询",
})
return cards
def _plans_from_evidence(cards: list[dict[str, str]], answer: str) -> list[dict[str, Any]]:
plans: list[dict[str, Any]] = []
for idx, card in enumerate(cards[:4], start=1):
plans.append({
"rank": idx,
"label": "图查询证据",
"plan_name": card.get("name") or f"图谱证据 {idx}",
"product_name": card.get("name") or "",
"fit_score": max(60, 92 - idx * 4),
"match_reasons": ["LLM 生成 Cypher 查询命中", "来自百姓惠知识图谱"],
"route_summary": card.get("summary") or answer[:240],
"quote_summary": answer[:360],
"variant_summary": card.get("summary") or "",
"daily_itinerary": [],
"hotels": [],
"restaurants": [],
"vehicles": [],
"policies": [],
"cost_breakdown": [],
"plan_kind": "llm_graph_evidence",
})
return plans
def _safe_confidence(value: Any, default: float = 0.72) -> float:
try:
return max(0.0, min(1.0, float(value)))
except Exception:
return default
async def answer_graph_question(
question: str,
graph_name: str,
*,
customer_context: dict[str, Any] | None = None,
limit: int = 80,
) -> dict[str, Any]:
started_at = time.perf_counter()
limit = min(max(int(limit or 80), 1), GRAPH_QA_MAX_LIMIT)
question = question.strip()
if not question:
raise ValueError("question required")
cache_key = _cache_key(question, graph_name, limit, customer_context)
cached_response = _get_cached_response(cache_key, started_at)
if cached_response:
return cached_response
intent_started = time.perf_counter()
rule_intent = _classify_graph_qa_intent(question)
intent, planner_trace = await _llm_plan_intent(question, graph_name, rule_intent, customer_context)
intent["_planner_trace"] = planner_trace
intent_ms = max(1, round((time.perf_counter() - intent_started) * 1000))
fast_response = await _fast_graph_response(
question,
graph_name,
limit=limit,
customer_context=customer_context,
started_at=started_at,
intent=intent,
intent_ms=intent_ms,
)
if fast_response:
_set_cached_response(cache_key, fast_response)
return fast_response
context_response = await _deterministic_context_response(
question,
graph_name,
limit=limit,
customer_context=customer_context,
started_at=started_at,
intent=intent,
intent_ms=intent_ms,
)
if context_response:
_set_cached_response(cache_key, context_response)
return context_response
cypher_cache_key = _cypher_cache_key(question, graph_name, limit, intent)
cached_decision = _get_cached_cypher(cypher_cache_key)
cypher_cache_hit = bool(cached_decision)
cypher_client: LlmClient | None = None
if cached_decision:
decision = cached_decision
schema_ms = 0
cypher_ms = 0
else:
cypher_client = await _graph_qa_client(max_tokens=1000, use_config_max_tokens=False)
if cypher_client is None:
raise RuntimeError("LLM 未配置,无法执行自然语言图查询")
cypher_client.timeout = min(float(getattr(cypher_client, "timeout", GRAPH_QA_CYPHER_TIMEOUT_SECONDS) or GRAPH_QA_CYPHER_TIMEOUT_SECONDS), GRAPH_QA_CYPHER_TIMEOUT_SECONDS)
schema_started = time.perf_counter()
schema = await asyncio.to_thread(_schema_snapshot, graph_name)
schema_ms = max(1, round((time.perf_counter() - schema_started) * 1000))
plan_payload = {
"question": question,
"graph_name": graph_name,
"intent": intent,
"customer_context": customer_context or {},
"graph_schema": schema,
"default_limit": limit,
}
cypher_started = time.perf_counter()
decision = await _chat_json_timed(
cypher_client,
GRAPH_QA_CYPHER_SYS,
json.dumps(plan_payload, ensure_ascii=False),
timeout_seconds=GRAPH_QA_CYPHER_TIMEOUT_SECONDS,
attempts=1,
)
cypher_ms = max(1, round((time.perf_counter() - cypher_started) * 1000))
_set_cached_cypher(cypher_cache_key, decision)
answer_client = await _graph_qa_client(max_tokens=650)
if answer_client is None:
raise RuntimeError("LLM 未配置,无法执行自然语言图查询")
answer_client.timeout = min(float(getattr(answer_client, "timeout", GRAPH_QA_ANSWER_TIMEOUT_SECONDS) or GRAPH_QA_ANSWER_TIMEOUT_SECONDS), GRAPH_QA_ANSWER_TIMEOUT_SECONDS)
generated_cypher = _clean_cypher(decision.get("cypher"), limit)
llm_generated_cypher = generated_cypher
graph_result: dict[str, Any]
query_error = ""
repaired = False
fallback_used = False
fallback_cypher = ""
fallback_ms = 0
try:
graph_started = time.perf_counter()
graph_result = await asyncio.to_thread(_run_cypher, graph_name, generated_cypher, limit)
graph_ms = max(1, round((time.perf_counter() - graph_started) * 1000))
except Exception as exc: # noqa: BLE001
query_error = str(exc)[:400]
if cypher_client is None:
cypher_client = await _graph_qa_client(max_tokens=1000, use_config_max_tokens=False)
if cypher_client is None:
raise
cypher_client.timeout = min(float(getattr(cypher_client, "timeout", GRAPH_QA_REPAIR_TIMEOUT_SECONDS) or GRAPH_QA_REPAIR_TIMEOUT_SECONDS), GRAPH_QA_REPAIR_TIMEOUT_SECONDS)
repair_payload = {
"question": question,
"graph_name": graph_name,
"intent": intent,
"customer_context": customer_context or {},
"failed_cypher": generated_cypher,
"execution_error": query_error,
}
repaired_decision = await _chat_json_timed(
cypher_client,
GRAPH_QA_REPAIR_SYS,
json.dumps(repair_payload, ensure_ascii=False),
timeout_seconds=GRAPH_QA_REPAIR_TIMEOUT_SECONDS,
attempts=1,
)
generated_cypher = _clean_cypher(repaired_decision.get("cypher"), limit)
decision = {**decision, **repaired_decision}
repaired = True
graph_started = time.perf_counter()
graph_result = await asyncio.to_thread(_run_cypher, graph_name, generated_cypher, limit)
graph_ms = max(1, round((time.perf_counter() - graph_started) * 1000))
if _is_fee_question(question) and not _has_fee_evidence(graph_result):
fallback_cypher = _fee_fallback_cypher(question, limit) or ""
if fallback_cypher:
try:
fallback_started = time.perf_counter()
fallback_result = await asyncio.to_thread(_run_cypher, graph_name, fallback_cypher, limit)
fallback_ms = max(1, round((time.perf_counter() - fallback_started) * 1000))
if fallback_result.get("row_count") and _has_fee_evidence(fallback_result):
graph_result = fallback_result
generated_cypher = fallback_cypher
fallback_used = True
except Exception as exc: # noqa: BLE001
query_error = "; ".join(item for item in [query_error, f"fallback: {str(exc)[:220]}"] if item)
if not graph_result.get("row_count"):
latency_ms = max(1, round((time.perf_counter() - started_at) * 1000))
planner = intent.get("_planner_trace") if isinstance(intent.get("_planner_trace"), dict) else {}
planner_ms = _safe_int(planner.get("latency_ms"), 0)
answer = "已完成 LLM 图查询,但当前图谱没有命中可确认数据。建议补充更具体的线路、景区、套餐名称、出发日期或人数后再查。"
customer_reply = "抱歉,当前图谱暂未查到相关产品或套餐。您可以补充具体线路、景区、出发日期或套餐名称,我再帮您核实。"
response = {
"question": question,
"graph_name": graph_name,
"answer": answer,
"customer_reply": customer_reply,
"copy_text": customer_reply,
"plans": [],
"evidence": [],
"sales_scripts": [],
"follow_up_questions": ["请补充具体线路/产品名称。", "请说明出发日期、人数或套餐名称。", "是否要改查百姓惠旅游线路、价格、景区或酒店?"],
"risk_notes": ["图谱查询结果为空,不代表业务一定不存在,需按产品库或人工渠道二次核实。"],
"confidence": 0.28,
"graph_result": graph_result,
"trace": {
"method": "llm_to_cypher_graph_qa_v1",
"query_source": "llm_generated_cypher",
"rule_query_used": False,
"llm_used": True,
"llm_error": "",
"generated_cypher": llm_generated_cypher,
"effective_cypher": generated_cypher,
"cache": {"response": False, "cypher": cypher_cache_hit},
"cypher_cache_hit": cypher_cache_hit,
"fallback_cypher": fallback_cypher,
"fallback_query_used": fallback_used,
"cypher_reason": decision.get("reason") or "",
"intent": intent.get("intent") or "",
"llm_planner_used": bool(planner.get("used")),
"llm_planner": planner,
"graph_qa_intent": intent,
"routing_strategy": "intent_template_first_llm_fallback",
"cypher_repaired": repaired,
"first_query_error": query_error,
"row_count": 0,
"node_count": 0,
"relationship_count": 0,
"latency_ms": latency_ms,
"stage_timings_ms": {
"intent_classification": intent_ms,
"llm_planner": planner_ms,
"schema": schema_ms,
"cypher_generation": cypher_ms,
"graph_query": graph_ms,
"fallback_graph_query": fallback_ms,
"answer_synthesis": 0,
},
"performance_target_ms": 1200,
"response_mode": "llm_graph_qa",
"response_mode_label": "LLM 图查询问答",
"graph_capabilities_used": ["LLM-to-Cypher", "FalkorDB 只读查询", "空结果快速返回"],
"retrieval_summary": {
"cypher": generated_cypher,
"llm_generated_cypher": llm_generated_cypher,
"fallback_query_used": fallback_used,
"rows": 0,
"nodes": 0,
"relationships": 0,
},
},
}
_set_cached_response(cache_key, response)
return response
answer_payload = {
"question": question,
"graph_name": graph_name,
"intent": intent,
"answer_focus": decision.get("answer_focus") or "",
"cypher": generated_cypher,
"graph_result": {
"columns": graph_result["columns"],
"row_count": graph_result["row_count"],
"rows": graph_result["rows"][:8],
"nodes": graph_result["nodes"][:10],
"relationships": graph_result["relationships"][:10],
},
"customer_context": customer_context or {},
}
answer_started = time.perf_counter()
answer_error = ""
try:
answer_data = await _chat_json_timed(
answer_client,
GRAPH_QA_ANSWER_SYS,
json.dumps(answer_payload, ensure_ascii=False),
timeout_seconds=GRAPH_QA_ANSWER_TIMEOUT_SECONDS,
attempts=1,
)
except asyncio.TimeoutError:
answer_data = {}
answer_error = "answer_timeout"
except Exception as exc: # noqa: BLE001
answer_data = {}
answer_error = str(exc)[:260]
answer_ms = max(1, round((time.perf_counter() - answer_started) * 1000))
answer = str(answer_data.get("answer") or answer_data.get("customer_reply") or "").strip()
route_list_override = _route_list_answer(question, graph_result)
if route_list_override:
answer, customer_reply_override = route_list_override
answer_data["customer_reply"] = customer_reply_override
answer_data.setdefault("follow_up_questions", ["要查哪条线路的价格?", "请提供出发日期和人数。"])
answer_data.setdefault("risk_notes", ["价格、余位和用车需按具体团期二次核实。"])
if not answer:
answer = "图查询已完成,但答案组织 LLM 未在时间预算内返回。为避免编造,建议查看返回的 knowledge.evidence,并补充更明确的景点、线路、团期、人数或费用口径。"
customer_reply = str(answer_data.get("customer_reply") or answer).strip()
evidence = _evidence_cards(answer_data, graph_result)
plans = _plans_from_evidence(evidence, answer)
latency_ms = max(1, round((time.perf_counter() - started_at) * 1000))
planner = intent.get("_planner_trace") if isinstance(intent.get("_planner_trace"), dict) else {}
planner_ms = _safe_int(planner.get("latency_ms"), 0)
response = {
"question": question,
"graph_name": graph_name,
"answer": answer,
"customer_reply": customer_reply,
"copy_text": customer_reply or answer,
"plans": plans,
"evidence": evidence,
"sales_scripts": [],
"follow_up_questions": _list_texts(answer_data.get("follow_up_questions"), 4, 120),
"risk_notes": _list_texts(answer_data.get("risk_notes"), 4, 160),
"confidence": _safe_confidence(answer_data.get("confidence")),
"graph_result": graph_result,
"trace": {
"method": "llm_to_cypher_graph_qa_v1",
"query_source": "llm_generated_cypher",
"rule_query_used": False,
"llm_used": True,
"llm_error": answer_error,
"generated_cypher": llm_generated_cypher,
"effective_cypher": generated_cypher,
"cache": {"response": False, "cypher": cypher_cache_hit},
"cypher_cache_hit": cypher_cache_hit,
"fallback_cypher": fallback_cypher,
"fallback_query_used": fallback_used,
"cypher_reason": decision.get("reason") or "",
"intent": intent.get("intent") or "",
"llm_planner_used": bool(planner.get("used")),
"llm_planner": planner,
"graph_qa_intent": intent,
"routing_strategy": "intent_template_first_llm_fallback",
"cypher_repaired": repaired,
"first_query_error": query_error,
"row_count": graph_result["row_count"],
"node_count": len(graph_result["nodes"]),
"relationship_count": len(graph_result["relationships"]),
"latency_ms": latency_ms,
"stage_timings_ms": {
"intent_classification": intent_ms,
"llm_planner": planner_ms,
"schema": schema_ms,
"cypher_generation": cypher_ms,
"graph_query": graph_ms,
"fallback_graph_query": fallback_ms,
"answer_synthesis": answer_ms,
},
"performance_target_ms": 1200,
"response_mode": "llm_graph_qa",
"response_mode_label": "LLM 图查询问答",
"graph_capabilities_used": ["LLM-to-Cypher", "FalkorDB 只读查询", "LLM evidence synthesis"],
"retrieval_summary": {
"cypher": generated_cypher,
"llm_generated_cypher": llm_generated_cypher,
"fallback_query_used": fallback_used,
"rows": graph_result["row_count"],
"nodes": len(graph_result["nodes"]),
"relationships": len(graph_result["relationships"]),
},
},
}
_set_cached_response(cache_key, response)
return response