vent_tools.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335
  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. """
  18. import json
  19. from fastmcp import Client
  20. from tools.mcp_logger import log_mcp_call
  21. # ============================================================
  22. # MCP 客户端辅助函数
  23. # ============================================================
  24. def _get_mcp_url() -> str:
  25. """获取 MCP 服务 URL。
  26. 从 src/.env 读取 MCP_BASE_URL 配置,拼接 /mcp 路径。
  27. """
  28. from dotenv import dotenv_values
  29. from pathlib import Path as _Path
  30. _env_file = _Path(__file__).parent.parent / ".env" # src/.env
  31. _cfg = dotenv_values(str(_env_file))
  32. mcp_url = _cfg.get("MCP_BASE_URL", "http://localhost:8100").rstrip("/")
  33. return f"{mcp_url}/mcp"
  34. async def _call_mcp_tool(tool_name: str, arguments: dict) -> str:
  35. """通过 FastMCP 异步客户端调用远程 MCP 工具。
  36. 参考 xfl_demo_client.py 的实现方式,使用 fastmcp.Client 异步上下文管理器。
  37. """
  38. mcp_url = _get_mcp_url()
  39. try:
  40. client = Client(mcp_url)
  41. async with client:
  42. result = await client.call_tool(tool_name, arguments)
  43. raw = result.content[0].text
  44. log_mcp_call(tool_name, arguments, raw)
  45. return raw
  46. except Exception as e:
  47. log_mcp_call(tool_name, arguments, "", error=str(e))
  48. return json.dumps({"error": f"MCP调用失败: {e}", "tool": tool_name}, ensure_ascii=False)
  49. # ============================================================
  50. # DeepAgent 工具函数(供 Agent 调用)
  51. # ============================================================
  52. async def query_device_data(device_id: str = "") -> str:
  53. """查询设备实时监测数据。
  54. 当需要获取设备实时数据时调用此工具。
  55. Args:
  56. device_id: 设备ID。
  57. Returns:
  58. JSON格式的监测数据。
  59. """
  60. return await _call_mcp_tool("query_device_data", {"device_id": device_id})
  61. async def query_device_data_by_id(device_id: str = "", page_size: int = 20) -> str:
  62. """根据设备ID查询实时数据和报警数据。
  63. 当需要获取某设备的实时监测数据和报警信息时调用此工具。
  64. Args:
  65. device_id: 设备ID。
  66. page_size: 返回数据条数,默认20。
  67. Returns:
  68. JSON格式的监测数据和报警数据。
  69. """
  70. return await _call_mcp_tool("query_device_data_by_id", {"device_id": device_id, "page_size": page_size})
  71. async def query_devices_by_tunnel(tunnel_name: str = "", device_type: str = None) -> str:
  72. """根据巷道名称查询该巷道绑定的设备及实时数据。
  73. 当需要查询某巷道下所有关联设备及其当前监测数据时调用此工具。
  74. Args:
  75. tunnel_name: 巷道名称(必填)。
  76. device_type: 设备类型,不传查询全部设备。
  77. Returns:
  78. JSON格式的设备列表及实时数据。
  79. """
  80. params = {"tunnel_name": tunnel_name}
  81. if device_type:
  82. params["device_type"] = device_type
  83. return await _call_mcp_tool("query_devices_by_tunnel", params)
  84. async def query_devices_by_tunnel_id(tunnel_id: str = "", model_id: str = None, device_type: str = None) -> str:
  85. """根据巷道ID查询已绑定设备。
  86. 当需要通过巷道ID查询该巷道下绑定的设备列表时调用此工具。
  87. Args:
  88. tunnel_id: 巷道ID(必填)。
  89. model_id: 模型ID,可选。
  90. device_type: 设备类型,可选。
  91. Returns:
  92. JSON格式的设备列表。
  93. """
  94. params = {"tunnel_id": tunnel_id}
  95. if model_id:
  96. params["model_id"] = model_id
  97. if device_type:
  98. params["device_type"] = device_type
  99. return await _call_mcp_tool("query_devices_by_tunnel_id", params)
  100. async def get_tun_list_by_modelid(model_id: str = "") -> str:
  101. """通过模型ID获取巷道列表。
  102. 当需要获取指定模型下的所有巷道列表时调用此工具。
  103. Args:
  104. model_id: 模型ID(必填)。
  105. Returns:
  106. JSON格式的巷道列表。
  107. """
  108. return await _call_mcp_tool("get_tun_list_by_modelid", {"model_id": model_id})
  109. async def query_tunnels_by_model(model_id: str = "") -> str:
  110. """根据模型ID查询巷道列表。
  111. 通过模型ID获取该模型下所有巷道的基本信息列表,
  112. 包括巷道ID、名称、类型、需风量等。
  113. Args:
  114. model_id: 模型ID(必填)。
  115. Returns:
  116. JSON格式的巷道列表。
  117. """
  118. return await _call_mcp_tool("query_tunnels_by_model", {"model_id": model_id})
  119. async def query_tunnel_list(model_id: int, tunnel_name: str) -> str:
  120. """根据模型ID和巷道名称(模糊)查询巷道列表。
  121. 当用户提到具体巷道名称时,**优先使用此工具**进行模糊匹配,
  122. 可大幅减少返回数据量。仅返回 modelId / tunnelId / tunnelName 三个字段。
  123. 若此工具无匹配结果,再回退使用 query_tunnels_by_model 获取全量列表。
  124. Args:
  125. model_id: 模型ID(必填),默认使用 .env 中的 DEFAULT_MODEL_ID。
  126. tunnel_name: 巷道名称(模糊匹配,必填)。
  127. Returns:
  128. JSON格式的巷道列表(仅含 modelId/tunnelId/tunnelName)。
  129. 包含多跳同名的巷道,因为一条完整的巷道是多节的,所以获取挂载设备时要多次判断。
  130. """
  131. return await _call_mcp_tool("query_tunnel_list", {
  132. "model_id": model_id,
  133. "tunnel_name": tunnel_name,
  134. })
  135. async def query_tun_data_by_id(tun_id: str = "") -> str:
  136. """根据ID查询指定巷道解读数据。
  137. 当需要获取某条巷道的风速、风量、瓦斯浓度、温度、设备状态等实时数据时调用此工具。
  138. Args:
  139. tun_id: 巷道ID。
  140. Returns:
  141. JSON格式的监测数据,包含风速、风量、瓦斯、温度、设备状态等信息。
  142. tunId 巷道唯一 ID
  143. tunnelName 巷道名称
  144. needAirVolume 巷道需配风量,单位 m³/min
  145. usingType 巷道类型 0-回采工作面;1-掘进工作面;2-辅运巷;3-主运巷;4-硐室;5-联络巷;6-进风井;7-回风井;8-专用回风巷
  146. usingTypeName 巷道用途中文名称(如掘进工作面)
  147. regulationId 风速规范配置 ID
  148. permissibleMin 允许最小风速,m/s
  149. permissibleMax 允许最大风速,m/s
  150. permissibleVelocity 巷道风速规范限制对象
  151. sensorIds 绑定的所有传感器 ID 数组
  152. deviceCount 绑定设备总数量
  153. devices 巷道绑定设备列表数组
  154. permissibleVelocity
  155. model_reg_id 风速规范主键 ID,同外层 regulationId
  156. fmin 最低允许风速 m/s
  157. fmax 最高允许风速 m/s
  158. sourceType 设备数据源类型:wind = 测风设备,model_sensor = 模型计算传感器
  159. sourceId 设备数据源唯一编号
  160. sourceName 数据源名称
  161. installPos 设备井下安装位置
  162. sensorId 传感器唯一标识 ID
  163. parentId 父级传感器 ID,单设备时与 sensorId 一致
  164. airVolume 实时风量,单位 m³/min;null 表示无实时数据
  165. windSpeed 实时风速,单位 m/s;null 表示无实时数据
  166. warnFlag 报警标识:0 = 无报警,非 0 代表存在异常报警
  167. netStatus 设备网络在线状态:1 = 在线,null / 其他 = 离线 / 无数据
  168. deviceStatus 设备运行状态编码
  169. deviceStatusName 设备运行状态中文描述(在线 / 离线等)
  170. deviceName 设备展示名称
  171. readTime 最新数据采集时间,格式 yyyy-MM-dd HH:mm:ss
  172. deviceType 设备类型标识,modelsensor_speed = 风速模型传感器
  173. readData 设备实时采集原始数据对象
  174. alarmDescription 当前设备单条报警文本,无报警为空
  175. alarmDescriptions 设备多条报警信息数组,无报警为空数组
  176. m3 实时风量数值 m³/min
  177. sign 风向标识
  178. tTime 传感器原始采集时间
  179. va 实时风速数值 m/s
  180. isRun 设备运行状态标记
  181. """
  182. return await _call_mcp_tool("query_wind_by_tunid", {"tun_id": tun_id, "model_id":2012326636757958658})
  183. def query_knowledge_base(question: str = "", top_k: int = 5) -> str:
  184. """查询煤矿安全知识库,检索《煤矿安全规程》条款依据、标准规范原文。
  185. 当需要查询某项安全条款的具体依据、规程出处、标准参数限值、
  186. 技术规范原文时调用此工具。知识库地址 http://39.97.59.228:8067,
  187. 使用 /api/retrieve 接口进行语义+关键词混合检索。
  188. Args:
  189. question: 检索问句,如"煤矿安全规程 瓦斯浓度限值 采煤工作面"、"AQ1056 风量计算方法"
  190. top_k: 返回结果数量,1~20,默认5
  191. Returns:
  192. JSON格式的检索结果,包含 context(相关段落文本)和 sources(来源列表)。
  193. """
  194. import httpx
  195. try:
  196. payload = {
  197. "question": question,
  198. "top_k": max(1, min(20, top_k)),
  199. "search_mode": "hybrid",
  200. "format": "compact",
  201. }
  202. with httpx.Client(timeout=15.0) as client:
  203. resp = client.post("http://39.97.59.228:8067/api/retrieve", json=payload)
  204. resp.raise_for_status()
  205. result = resp.json()
  206. print(result.get("context", "")[:200])
  207. print(json.dumps(result, indent=2, ensure_ascii=False))
  208. return json.dumps(result, ensure_ascii=False)
  209. except Exception as e:
  210. return json.dumps({
  211. "error": f"知识库查询失败: {e}",
  212. "question": question,
  213. }, ensure_ascii=False)
  214. # ============================================================
  215. # 需风量数据查询(通防管控平台 MCP)
  216. # ============================================================
  217. async def get_needq_all_data() -> str:
  218. """从通防管控平台获取全部用风地点的需风量数据。
  219. 调用远程 MCP 工具 get_needq_all_data,返回管控平台上所有用风地点的
  220. 计划需风量、实际风量、偏差等数据。用于回答用户"查看管控平台的需风量情况"
  221. "通防管控平台上各地点需风量是多少"等问题。
  222. Returns:
  223. JSON 格式的全部需风量数据,包含各用风地点的名称、类型、计划需风量、
  224. 实际风量、偏差等字段。
  225. """
  226. return await _call_mcp_tool("get_needq_all_data", {})
  227. # ============================================================
  228. # 监测历史数据查询(通防管控平台 MCP)
  229. # ============================================================
  230. async def list_ventanaly_monitor_data_days(
  231. strtype: str = "",
  232. gdeviceids: str = "",
  233. ttime_begin: str = "",
  234. ttime_end: str = "",
  235. device_num: str = "",
  236. skip: int = 8,
  237. page_no: int = 1,
  238. page_size: int = 100,
  239. # column: str = "",
  240. ) -> str:
  241. """查询监测设备历史时序数据。
  242. 根据设备类型、设备ID、时间范围等条件,查询监测设备的历史时序数据。
  243. 适用于查询风速、风量、瓦斯、温度等传感器在一段时间内的历史趋势。
  244. Args:
  245. strtype: 设备类型(必填),对应后端 strtype,如 "fanmain_stem_wp_2"
  246. gdeviceids: 设备ID(必填),如 "11111004";多个设备按后端要求格式拼接
  247. ttime_begin: 开始时间(必填),格式 "yyyy-MM-dd HH:mm:ss"
  248. ttime_end: 结束时间(必填),格式 "yyyy-MM-dd HH:mm:ss"
  249. device_num: 设备编号(必填),对应后端 deviceNum,如 "Fan1"
  250. skip: 采样间隔/跳点参数(必填),默认8,1=5秒//2=10秒//3=30秒//4=1分钟//5=5分钟//6=10分钟//7=30分钟8=1小时,如果半天内默认按10分钟查询,如果需要跨天则默认按1小时查询,其他情况请按合适的采样间隔获取,考虑数据库的压力。
  251. page_no: 页码,默认 1
  252. page_size: 每页条数,默认 200
  253. column: 排序字段,可选
  254. Returns:
  255. JSON格式的历史监测时序数据。
  256. """
  257. return await _call_mcp_tool("list_ventanaly_monitor_data_days", {
  258. "strtype": strtype,
  259. "gdeviceids": gdeviceids,
  260. "ttime_begin": ttime_begin,
  261. "ttime_end": ttime_end,
  262. "device_num": device_num,
  263. "skip": skip,
  264. "page_no": page_no,
  265. "page_size": page_size,
  266. # "column": column,
  267. })