vent_tools.py 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509
  1. # -*- coding: utf-8 -*-
  2. """
  3. DeepAgents 工具函数模块
  4. 所有工具函数供 DeepAgent 调用,通过 MCP 客户端获取模拟数据。
  5. 工具函数返回结构化 JSON 文本,由 LLM 解析并生成自然语言回复。
  6. 工具列表:
  7. - query_device_data: 查询设备实时数据
  8. - query_device_data_by_id: 根据设备ID查询实时数据和报警数据
  9. - query_devices_by_tunnel: 根据巷道名称查询绑定设备
  10. - query_devices_by_tunnel_id: 根据巷道ID查询绑定设备
  11. - get_tun_list_by_modelid: 通过模型ID获取巷道列表(旧接口)
  12. - query_tunnels_by_model: 根据模型ID查询巷道列表(新接口)
  13. - query_tun_data_by_id: 根据巷道ID查询实时监测数据
  14. - query_knowledge_base: 查询煤矿安全知识库
  15. - get_needq_all_data: 获取全部需风量数据
  16. - list_ventanaly_monitor_data_days: 查询监测历史时序数据
  17. - get_device_kind_dict: 查询设备大类小类全量字典
  18. - get_device_list_by_kind: 根据设备类型查询设备列表
  19. - query_device_realtime_data: 查询设备实时监测数据
  20. """
  21. import json
  22. import contextvars
  23. from fastmcp import Client
  24. from tools.mcp_logger import log_mcp_call
  25. # ============================================================
  26. # 当前用户上下文(供工具函数获取调用者身份)
  27. # ============================================================
  28. _current_user: contextvars.ContextVar[str] = contextvars.ContextVar(
  29. 'current_user', default='admin'
  30. )
  31. # ============================================================
  32. # MCP 客户端辅助函数
  33. # ============================================================
  34. def _get_mcp_url() -> str:
  35. """获取 MCP 服务 URL。
  36. 从 src/.env 读取 MCP_BASE_URL 配置,拼接 /mcp 路径。
  37. """
  38. from dotenv import dotenv_values
  39. from pathlib import Path as _Path
  40. _env_file = _Path(__file__).parent.parent / ".env" # src/.env
  41. _cfg = dotenv_values(str(_env_file))
  42. mcp_url = _cfg.get("MCP_BASE_URL", "http://localhost:8100").rstrip("/")
  43. return f"{mcp_url}/mcp"
  44. async def _call_mcp_tool(tool_name: str, arguments: dict) -> str:
  45. """通过 FastMCP 异步客户端调用远程 MCP 工具。
  46. 参考 xfl_demo_client.py 的实现方式,使用 fastmcp.Client 异步上下文管理器。
  47. """
  48. mcp_url = _get_mcp_url()
  49. try:
  50. client = Client(mcp_url)
  51. async with client:
  52. result = await client.call_tool(tool_name, arguments)
  53. raw = result.content[0].text
  54. log_mcp_call(tool_name, arguments, raw)
  55. return raw
  56. except Exception as e:
  57. log_mcp_call(tool_name, arguments, "", error=str(e))
  58. return json.dumps({"error": f"MCP调用失败: {e}", "tool": tool_name}, ensure_ascii=False)
  59. # ============================================================
  60. # DeepAgent 工具函数(供 Agent 调用)
  61. # ============================================================
  62. async def query_device_data(device_id: str = "") -> str:
  63. """查询设备实时监测数据。
  64. 当需要获取设备实时数据时调用此工具。
  65. Args:
  66. device_id: 设备ID。
  67. Returns:
  68. JSON格式的监测数据。
  69. """
  70. return await _call_mcp_tool("query_device_data", {"device_id": device_id})
  71. async def query_device_data_by_id(device_id: str = "", page_size: int = 20) -> str:
  72. """根据设备ID查询实时数据和报警数据。
  73. 当需要获取某设备的实时监测数据和报警信息时调用此工具。
  74. Args:
  75. device_id: 设备ID。
  76. page_size: 返回数据条数,默认20。
  77. Returns:
  78. JSON格式的监测数据和报警数据。
  79. """
  80. return await _call_mcp_tool("query_device_data_by_id", {"device_id": device_id, "page_size": page_size})
  81. async def query_devices_by_tunnel(tunnel_name: str = "", device_type: str = None) -> str:
  82. """根据巷道名称查询该巷道绑定的设备及实时数据。
  83. 当需要查询某巷道下所有关联设备及其当前监测数据时调用此工具。
  84. Args:
  85. tunnel_name: 巷道名称(必填)。
  86. device_type: 设备类型,不传查询全部设备。
  87. Returns:
  88. JSON格式的设备列表及实时数据。
  89. """
  90. params = {"tunnel_name": tunnel_name}
  91. if device_type:
  92. params["device_type"] = device_type
  93. return await _call_mcp_tool("query_devices_by_tunnel", params)
  94. async def query_devices_by_tunnel_id(tunnel_id: str = "", model_id: str = None, device_type: str = None) -> str:
  95. """根据巷道ID查询已绑定设备。
  96. 当需要通过巷道ID查询该巷道下绑定的设备列表时调用此工具。
  97. Args:
  98. tunnel_id: 巷道ID(必填)。
  99. model_id: 模型ID,可选。
  100. device_type: 设备类型,可选。
  101. Returns:
  102. JSON格式的设备列表。
  103. """
  104. params = {"tunnel_id": tunnel_id}
  105. if model_id:
  106. params["model_id"] = model_id
  107. if device_type:
  108. params["device_type"] = device_type
  109. return await _call_mcp_tool("query_devices_by_tunnel_id", params)
  110. async def get_tun_list_by_modelid(model_id: str = "") -> str:
  111. """通过模型ID获取巷道列表。
  112. 当需要获取指定模型下的所有巷道列表时调用此工具。
  113. Args:
  114. model_id: 模型ID(必填)。
  115. Returns:
  116. JSON格式的巷道列表。
  117. """
  118. return await _call_mcp_tool("get_tun_list_by_modelid", {"model_id": model_id})
  119. async def query_tunnels_by_model(model_id: str = "") -> str:
  120. """根据模型ID查询巷道列表。
  121. 通过模型ID获取该模型下所有巷道的基本信息列表,
  122. 包括巷道ID、名称、类型、需风量等。
  123. Args:
  124. model_id: 模型ID(必填)。
  125. Returns:
  126. JSON格式的巷道列表。
  127. """
  128. return await _call_mcp_tool("query_tunnels_by_model", {"model_id": model_id})
  129. async def query_tunnel_list(model_id: int, tunnel_name: str) -> str:
  130. """根据模型ID和巷道名称(模糊)查询巷道列表。
  131. 当用户提到具体巷道名称时,**优先使用此工具**进行模糊匹配,
  132. 可大幅减少返回数据量。仅返回 modelId / tunnelId / tunnelName 三个字段。
  133. 若此工具无匹配结果,再回退使用 query_tunnels_by_model 获取全量列表。
  134. Args:
  135. model_id: 模型ID(必填),默认使用 .env 中的 DEFAULT_MODEL_ID。
  136. tunnel_name: 巷道名称(模糊匹配,必填)。
  137. Returns:
  138. JSON格式的巷道列表(仅含 modelId/tunnelId/tunnelName)。
  139. 包含多跳同名的巷道,因为一条完整的巷道是多节的,所以获取挂载设备时要多次判断。
  140. """
  141. return await _call_mcp_tool("query_tunnel_list", {
  142. "model_id": model_id,
  143. "tunnel_name": tunnel_name,
  144. })
  145. async def query_tun_data_by_id(tun_id: str = "") -> str:
  146. """根据ID查询指定巷道解读数据。
  147. 当需要获取某条巷道的风速、风量、瓦斯浓度、温度、设备状态等实时数据时调用此工具。
  148. Args:
  149. tun_id: 巷道ID。
  150. Returns:
  151. JSON格式的监测数据,包含风速、风量、瓦斯、温度、设备状态等信息。
  152. tunId 巷道唯一 ID
  153. tunnelName 巷道名称
  154. needAirVolume 巷道需配风量,单位 m³/min
  155. usingType 巷道类型 0-回采工作面;1-掘进工作面;2-辅运巷;3-主运巷;4-硐室;5-联络巷;6-进风井;7-回风井;8-专用回风巷
  156. usingTypeName 巷道用途中文名称(如掘进工作面)
  157. regulationId 风速规范配置 ID
  158. permissibleMin 允许最小风速,m/s
  159. permissibleMax 允许最大风速,m/s
  160. permissibleVelocity 巷道风速规范限制对象
  161. sensorIds 绑定的所有传感器 ID 数组
  162. deviceCount 绑定设备总数量
  163. devices 巷道绑定设备列表数组
  164. permissibleVelocity
  165. model_reg_id 风速规范主键 ID,同外层 regulationId
  166. fmin 最低允许风速 m/s
  167. fmax 最高允许风速 m/s
  168. sourceType 设备数据源类型:wind = 测风设备,model_sensor = 模型计算传感器
  169. sourceId 设备数据源唯一编号
  170. sourceName 数据源名称
  171. installPos 设备井下安装位置
  172. sensorId 传感器唯一标识 ID
  173. parentId 父级传感器 ID,单设备时与 sensorId 一致
  174. airVolume 实时风量,单位 m³/min;null 表示无实时数据
  175. windSpeed 实时风速,单位 m/s;null 表示无实时数据
  176. warnFlag 报警标识:0 = 无报警,非 0 代表存在异常报警
  177. netStatus 设备网络在线状态:1 = 在线,null / 其他 = 离线 / 无数据
  178. deviceStatus 设备运行状态编码
  179. deviceStatusName 设备运行状态中文描述(在线 / 离线等)
  180. deviceName 设备展示名称
  181. readTime 最新数据采集时间,格式 yyyy-MM-dd HH:mm:ss
  182. deviceType 设备类型标识,modelsensor_speed = 风速模型传感器
  183. readData 设备实时采集原始数据对象
  184. alarmDescription 当前设备单条报警文本,无报警为空
  185. alarmDescriptions 设备多条报警信息数组,无报警为空数组
  186. m3 实时风量数值 m³/min
  187. sign 风向标识
  188. tTime 传感器原始采集时间
  189. va 实时风速数值 m/s
  190. isRun 设备运行状态标记
  191. """
  192. return await _call_mcp_tool("query_wind_by_tunid", {"tun_id": tun_id, "model_id":2012326636757958658})
  193. def query_knowledge_base(question: str = "", top_k: int = 5) -> str:
  194. """查询煤矿安全知识库,检索《煤矿安全规程》条款依据、标准规范原文。
  195. 当需要查询某项安全条款的具体依据、规程出处、标准参数限值、
  196. 技术规范原文时调用此工具。知识库地址 http://39.97.59.228:8067,
  197. 使用 /api/retrieve 接口进行语义+关键词混合检索。
  198. Args:
  199. question: 检索问句,如"煤矿安全规程 瓦斯浓度限值 采煤工作面"、"AQ1056 风量计算方法"
  200. top_k: 返回结果数量,1~20,默认5
  201. Returns:
  202. JSON格式的检索结果,包含 context(相关段落文本)和 sources(来源列表)。
  203. """
  204. import httpx
  205. try:
  206. payload = {
  207. "question": question,
  208. "top_k": max(1, min(20, top_k)),
  209. "search_mode": "hybrid",
  210. "format": "compact",
  211. }
  212. with httpx.Client(timeout=15.0) as client:
  213. resp = client.post("http://39.97.59.228:8067/api/retrieve", json=payload)
  214. resp.raise_for_status()
  215. result = resp.json()
  216. print(result.get("context", "")[:200])
  217. print(json.dumps(result, indent=2, ensure_ascii=False))
  218. return json.dumps(result, ensure_ascii=False)
  219. except Exception as e:
  220. return json.dumps({
  221. "error": f"知识库查询失败: {e}",
  222. "question": question,
  223. }, ensure_ascii=False)
  224. # ============================================================
  225. # 需风量数据查询(通防管控平台 MCP)
  226. # ============================================================
  227. async def get_needq_all_data() -> str:
  228. """从通防管控平台获取全部用风地点的需风量数据。
  229. 调用远程 MCP 工具 get_needq_all_data,返回管控平台上所有用风地点的
  230. 计划需风量、实际风量、偏差等数据。用于回答用户"查看管控平台的需风量情况"
  231. "通防管控平台上各地点需风量是多少"等问题。
  232. Returns:
  233. JSON 格式的全部需风量数据,包含各用风地点的名称、类型、计划需风量、
  234. 实际风量、偏差等字段。
  235. """
  236. return await _call_mcp_tool("get_needq_all_data", {})
  237. async def get_device_kind_dict() -> str:
  238. """查询设备大类(deviceKind)和小类(strType)的全量字典。
  239. 返回设备类型编码与中文名称的完整映射表,包含 deviceKind(设备大类)
  240. 和 strType(设备小类)两个维度的编码-名称对照。
  241. 无参数。
  242. 适用场景:
  243. - 用户询问"有哪些设备类型""设备分类有哪些""设备大类小类有哪些"
  244. - 需要将设备类型编码翻译为中文名称
  245. - 查询某种设备类型对应的 strType 编码以便后续调用历史数据工具
  246. - 用户想知道系统中支持监控哪些类型的设备
  247. Returns:
  248. JSON 格式的设备类型字典,包含 deviceKind 和 strType 映射。
  249. """
  250. return await _call_mcp_tool("get_device_kind_dict", {})
  251. async def get_device_list_by_kind(device_kind: str) -> str:
  252. """根据设备类型(deviceKind)查询该类型下的全量设备列表。
  253. 返回指定设备大类下的所有设备,包含设备ID、设备名称、安装位置、分站名称。
  254. 适用场景:
  255. - 用户询问"列出所有风速传感器""有哪些甲烷传感器"等按类型筛选设备
  256. - 需要获取某类设备的完整清单以便进一步查询实时/历史数据
  257. - 结合 get_device_kind_dict 先获取类型编码,再按类型查设备列表
  258. Args:
  259. device_kind: 设备大类编码(必填),如 "modelsensor_speed"表示风速传感器。
  260. 可先调用 get_device_kind_dict 获取所有可用的 deviceKind 编码。
  261. Returns:
  262. JSON 格式的设备列表,包含 device_id、device_name、install_pos、station_name。
  263. """
  264. return await _call_mcp_tool("get_device_list_by_kind", {"device_kind": device_kind})
  265. async def query_device_realtime_data(device_id: str) -> str:
  266. """查询指定设备的实时监测数据。
  267. 根据设备ID获取该设备当前最新的监测读数、运行状态和报警信息。
  268. 与 query_device_data_by_id 不同,本工具专注于单设备的实时快照数据,
  269. 返回结构更精简,延迟更低。
  270. 适用场景:
  271. - 用户询问"设备XXX的实时数据""传感器XXX当前读数是多少"
  272. - 已知设备ID,需要快速获取其最新监测值
  273. - 配合 get_device_list_by_kind 先查出设备ID列表,再逐个查询实时数据
  274. Args:
  275. device_id: 设备ID(必填),可从 query_devices_by_tunnel / get_device_list_by_kind 返回结果中获取。
  276. Returns:
  277. JSON 格式的设备实时监测数据,包含当前读数、单位、采集时间、在线状态、报警信息。
  278. """
  279. return await _call_mcp_tool("query_device_realtime_data", {"device_id": device_id})
  280. # ============================================================
  281. # 监测历史数据查询(通防管控平台 MCP)
  282. # ============================================================
  283. async def list_ventanaly_monitor_data_days(
  284. strtype: str = "",
  285. gdeviceids: str = "",
  286. ttime_begin: str = "",
  287. ttime_end: str = "",
  288. device_num: str = "",
  289. skip: int = 8,
  290. page_no: int = 1,
  291. page_size: int = 100,
  292. # column: str = "",
  293. ) -> str:
  294. """查询监测设备历史时序数据。
  295. 根据设备类型、设备ID、时间范围等条件,查询监测设备的历史时序数据。
  296. 适用于查询风速、风量、瓦斯、温度等传感器在一段时间内的历史趋势。
  297. Args:
  298. strtype: 设备类型(必填),对应后端 strtype,如 "fanmain_stem_wp_2"
  299. gdeviceids: 设备ID(必填),如 "11111004";多个设备按后端要求格式拼接
  300. ttime_begin: 开始时间(必填),格式 "yyyy-MM-dd HH:mm:ss"
  301. ttime_end: 结束时间(必填),格式 "yyyy-MM-dd HH:mm:ss"
  302. device_num: 设备编号(必填),对应后端 deviceNum,如 "Fan1"
  303. skip: 采样间隔/跳点参数(必填),默认8,1=5秒//2=10秒//3=30秒//4=1分钟//5=5分钟//6=10分钟//7=30分钟8=1小时,如果半天内默认按10分钟查询,如果需要跨天则默认按1小时查询,其他情况请按合适的采样间隔获取,考虑数据库的压力。
  304. page_no: 页码,默认 1
  305. page_size: 每页条数,默认 200
  306. column: 排序字段,可选
  307. Returns:
  308. JSON格式的历史监测时序数据。
  309. """
  310. return await _call_mcp_tool("list_ventanaly_monitor_data_days", {
  311. "strtype": strtype,
  312. "gdeviceids": gdeviceids,
  313. "ttime_begin": ttime_begin,
  314. "ttime_end": ttime_end,
  315. "device_num": device_num,
  316. "skip": skip,
  317. "page_no": page_no,
  318. "page_size": page_size,
  319. # "column": column,
  320. })
  321. # ============================================================
  322. # 用户偏好记忆工具
  323. # ============================================================
  324. def save_user_preference(content: str, keywords: str = "",
  325. category: str = "通用") -> str:
  326. """保存用户偏好/习惯到个人记忆库。
  327. 当用户在对话中明确要求"记住""保存为习惯/偏好""以后都用这个"时调用。
  328. 下次该用户对话时,系统会自动注入已保存的偏好作为上下文。
  329. Args:
  330. content: 偏好内容,描述具体的习惯或个性化要求。例如"习惯使用 m³/s 而非 m³/min"、"15216工作面默认采高3.5m"
  331. keywords: 触发关键词,多个用逗号分隔。例如"风速,单位,风量"。留空则自动匹配。
  332. category: 偏好分类,默认"通用"。可选值:需风量计算、数据解读、规程查询、通用
  333. Returns:
  334. 保存结果,含记录ID供后续删除用。
  335. """
  336. from db.chat_store import save_user_preference as _db_save
  337. user_name = _current_user.get()
  338. pref_id = _db_save(user_name, content, keywords, category)
  339. return json.dumps({
  340. "success": True,
  341. "id": pref_id,
  342. "message": f"已保存偏好 (id={pref_id}):{content}",
  343. }, ensure_ascii=False)
  344. def list_user_preferences() -> str:
  345. """查看当前用户已保存的所有偏好/习惯。
  346. 列出该用户所有偏好记录,包含ID、分类、内容、关键词、保存时间。
  347. Returns:
  348. JSON格式的偏好列表。
  349. """
  350. from db.chat_store import get_user_preferences as _db_list
  351. user_name = _current_user.get()
  352. prefs = _db_list(user_name)
  353. if not prefs:
  354. return json.dumps({"preferences": [], "message": "暂无保存的偏好"}, ensure_ascii=False)
  355. return json.dumps({
  356. "preferences": [
  357. {"id": p["id"], "category": p["category"], "content": p["content"],
  358. "keywords": p["keywords"], "created_at": p["created_at"]}
  359. for p in prefs
  360. ],
  361. }, ensure_ascii=False)
  362. def delete_user_preference(preference_id: int) -> str:
  363. """删除一条用户偏好记录。
  364. Args:
  365. preference_id: 要删除的偏好记录ID(可从 list_user_preferences 获取)
  366. Returns:
  367. 删除结果。
  368. """
  369. from db.chat_store import delete_user_preference as _db_delete
  370. user_name = _current_user.get()
  371. ok = _db_delete(preference_id, user_name)
  372. if ok:
  373. return json.dumps({"success": True, "message": f"已删除偏好 (id={preference_id})"}, ensure_ascii=False)
  374. return json.dumps({"success": False, "message": f"未找到偏好记录 id={preference_id} 或无权操作"}, ensure_ascii=False)
  375. # ============================================================
  376. # 计划审批工具(Human-in-the-Loop)
  377. # ============================================================
  378. def request_plan_approval(plan_summary: str) -> str:
  379. """提交执行计划等待人工审批。在制定好完整计划后调用此工具。
  380. 仅在计划模式(plan mode)下由 Agent 主动调用。
  381. 调用后会暂停执行,等待用户在前端审批(批准/拒绝)。
  382. 审批通过后自动继续执行计划。
  383. Args:
  384. plan_summary: 执行计划的简要描述,需包含:
  385. - 计划分几步,每步做什么
  386. - 每步预期调用哪些工具
  387. - 预期的输出结果
  388. Returns:
  389. "计划已批准,开始执行。" 或 "计划被拒绝。"
  390. """
  391. from langgraph.types import interrupt
  392. result = interrupt({
  393. "type": "plan_approval",
  394. "plan": plan_summary,
  395. "message": "智能体已制定执行计划,等待您的审批...",
  396. })
  397. if isinstance(result, dict) and result.get("action") == "approve":
  398. return "计划已批准,开始执行。"
  399. else:
  400. return "计划被拒绝。"