# -*- coding: utf-8 -*- """ DeepAgents 工具函数模块 所有工具函数供 DeepAgent 调用,通过 MCP 客户端获取模拟数据。 工具函数返回结构化 JSON 文本,由 LLM 解析并生成自然语言回复。 工具列表: - query_device_data: 查询设备实时数据 - query_device_data_by_id: 根据设备ID查询实时数据和报警数据 - query_devices_by_tunnel: 根据巷道名称查询绑定设备 - query_devices_by_tunnel_id: 根据巷道ID查询绑定设备 - get_tun_list_by_modelid: 通过模型ID获取巷道列表(旧接口) - query_tunnels_by_model: 根据模型ID查询巷道列表(新接口) - query_tun_data_by_id: 根据巷道ID查询实时监测数据 - query_knowledge_base: 查询煤矿安全知识库 - get_needq_all_data: 获取全部需风量数据 - list_ventanaly_monitor_data_days: 查询监测历史时序数据 """ import json from fastmcp import Client from tools.mcp_logger import log_mcp_call # ============================================================ # MCP 客户端辅助函数 # ============================================================ def _get_mcp_url() -> str: """获取 MCP 服务 URL。 从 src/.env 读取 MCP_BASE_URL 配置,拼接 /mcp 路径。 """ from dotenv import dotenv_values from pathlib import Path as _Path _env_file = _Path(__file__).parent.parent / ".env" # src/.env _cfg = dotenv_values(str(_env_file)) mcp_url = _cfg.get("MCP_BASE_URL", "http://localhost:8100").rstrip("/") return f"{mcp_url}/mcp" async def _call_mcp_tool(tool_name: str, arguments: dict) -> str: """通过 FastMCP 异步客户端调用远程 MCP 工具。 参考 xfl_demo_client.py 的实现方式,使用 fastmcp.Client 异步上下文管理器。 """ mcp_url = _get_mcp_url() try: client = Client(mcp_url) async with client: result = await client.call_tool(tool_name, arguments) raw = result.content[0].text log_mcp_call(tool_name, arguments, raw) return raw except Exception as e: log_mcp_call(tool_name, arguments, "", error=str(e)) return json.dumps({"error": f"MCP调用失败: {e}", "tool": tool_name}, ensure_ascii=False) # ============================================================ # DeepAgent 工具函数(供 Agent 调用) # ============================================================ async def query_device_data(device_id: str = "") -> str: """查询设备实时监测数据。 当需要获取设备实时数据时调用此工具。 Args: device_id: 设备ID。 Returns: JSON格式的监测数据。 """ return await _call_mcp_tool("query_device_data", {"device_id": device_id}) async def query_device_data_by_id(device_id: str = "", page_size: int = 20) -> str: """根据设备ID查询实时数据和报警数据。 当需要获取某设备的实时监测数据和报警信息时调用此工具。 Args: device_id: 设备ID。 page_size: 返回数据条数,默认20。 Returns: JSON格式的监测数据和报警数据。 """ return await _call_mcp_tool("query_device_data_by_id", {"device_id": device_id, "page_size": page_size}) async def query_devices_by_tunnel(tunnel_name: str = "", device_type: str = None) -> str: """根据巷道名称查询该巷道绑定的设备及实时数据。 当需要查询某巷道下所有关联设备及其当前监测数据时调用此工具。 Args: tunnel_name: 巷道名称(必填)。 device_type: 设备类型,不传查询全部设备。 Returns: JSON格式的设备列表及实时数据。 """ params = {"tunnel_name": tunnel_name} if device_type: params["device_type"] = device_type return await _call_mcp_tool("query_devices_by_tunnel", params) async def query_devices_by_tunnel_id(tunnel_id: str = "", model_id: str = None, device_type: str = None) -> str: """根据巷道ID查询已绑定设备。 当需要通过巷道ID查询该巷道下绑定的设备列表时调用此工具。 Args: tunnel_id: 巷道ID(必填)。 model_id: 模型ID,可选。 device_type: 设备类型,可选。 Returns: JSON格式的设备列表。 """ params = {"tunnel_id": tunnel_id} if model_id: params["model_id"] = model_id if device_type: params["device_type"] = device_type return await _call_mcp_tool("query_devices_by_tunnel_id", params) async def get_tun_list_by_modelid(model_id: str = "") -> str: """通过模型ID获取巷道列表。 当需要获取指定模型下的所有巷道列表时调用此工具。 Args: model_id: 模型ID(必填)。 Returns: JSON格式的巷道列表。 """ return await _call_mcp_tool("get_tun_list_by_modelid", {"model_id": model_id}) async def query_tunnels_by_model(model_id: str = "") -> str: """根据模型ID查询巷道列表。 通过模型ID获取该模型下所有巷道的基本信息列表, 包括巷道ID、名称、类型、需风量等。 Args: model_id: 模型ID(必填)。 Returns: JSON格式的巷道列表。 """ return await _call_mcp_tool("query_tunnels_by_model", {"model_id": model_id}) async def query_tunnel_list(model_id: int, tunnel_name: str) -> str: """根据模型ID和巷道名称(模糊)查询巷道列表。 当用户提到具体巷道名称时,**优先使用此工具**进行模糊匹配, 可大幅减少返回数据量。仅返回 modelId / tunnelId / tunnelName 三个字段。 若此工具无匹配结果,再回退使用 query_tunnels_by_model 获取全量列表。 Args: model_id: 模型ID(必填),默认使用 .env 中的 DEFAULT_MODEL_ID。 tunnel_name: 巷道名称(模糊匹配,必填)。 Returns: JSON格式的巷道列表(仅含 modelId/tunnelId/tunnelName)。 包含多跳同名的巷道,因为一条完整的巷道是多节的,所以获取挂载设备时要多次判断。 """ return await _call_mcp_tool("query_tunnel_list", { "model_id": model_id, "tunnel_name": tunnel_name, }) async def query_tun_data_by_id(tun_id: str = "") -> str: """根据ID查询指定巷道解读数据。 当需要获取某条巷道的风速、风量、瓦斯浓度、温度、设备状态等实时数据时调用此工具。 Args: tun_id: 巷道ID。 Returns: JSON格式的监测数据,包含风速、风量、瓦斯、温度、设备状态等信息。 tunId 巷道唯一 ID tunnelName 巷道名称 needAirVolume 巷道需配风量,单位 m³/min usingType 巷道类型 0-回采工作面;1-掘进工作面;2-辅运巷;3-主运巷;4-硐室;5-联络巷;6-进风井;7-回风井;8-专用回风巷 usingTypeName 巷道用途中文名称(如掘进工作面) regulationId 风速规范配置 ID permissibleMin 允许最小风速,m/s permissibleMax 允许最大风速,m/s permissibleVelocity 巷道风速规范限制对象 sensorIds 绑定的所有传感器 ID 数组 deviceCount 绑定设备总数量 devices 巷道绑定设备列表数组 permissibleVelocity model_reg_id 风速规范主键 ID,同外层 regulationId fmin 最低允许风速 m/s fmax 最高允许风速 m/s sourceType 设备数据源类型:wind = 测风设备,model_sensor = 模型计算传感器 sourceId 设备数据源唯一编号 sourceName 数据源名称 installPos 设备井下安装位置 sensorId 传感器唯一标识 ID parentId 父级传感器 ID,单设备时与 sensorId 一致 airVolume 实时风量,单位 m³/min;null 表示无实时数据 windSpeed 实时风速,单位 m/s;null 表示无实时数据 warnFlag 报警标识:0 = 无报警,非 0 代表存在异常报警 netStatus 设备网络在线状态:1 = 在线,null / 其他 = 离线 / 无数据 deviceStatus 设备运行状态编码 deviceStatusName 设备运行状态中文描述(在线 / 离线等) deviceName 设备展示名称 readTime 最新数据采集时间,格式 yyyy-MM-dd HH:mm:ss deviceType 设备类型标识,modelsensor_speed = 风速模型传感器 readData 设备实时采集原始数据对象 alarmDescription 当前设备单条报警文本,无报警为空 alarmDescriptions 设备多条报警信息数组,无报警为空数组 m3 实时风量数值 m³/min sign 风向标识 tTime 传感器原始采集时间 va 实时风速数值 m/s isRun 设备运行状态标记 """ return await _call_mcp_tool("query_wind_by_tunid", {"tun_id": tun_id, "model_id":2012326636757958658}) def query_knowledge_base(question: str = "", top_k: int = 5) -> str: """查询煤矿安全知识库,检索《煤矿安全规程》条款依据、标准规范原文。 当需要查询某项安全条款的具体依据、规程出处、标准参数限值、 技术规范原文时调用此工具。知识库地址 http://39.97.59.228:8067, 使用 /api/retrieve 接口进行语义+关键词混合检索。 Args: question: 检索问句,如"煤矿安全规程 瓦斯浓度限值 采煤工作面"、"AQ1056 风量计算方法" top_k: 返回结果数量,1~20,默认5 Returns: JSON格式的检索结果,包含 context(相关段落文本)和 sources(来源列表)。 """ import httpx try: payload = { "question": question, "top_k": max(1, min(20, top_k)), "search_mode": "hybrid", "format": "compact", } with httpx.Client(timeout=15.0) as client: resp = client.post("http://39.97.59.228:8067/api/retrieve", json=payload) resp.raise_for_status() result = resp.json() print(result.get("context", "")[:200]) print(json.dumps(result, indent=2, ensure_ascii=False)) return json.dumps(result, ensure_ascii=False) except Exception as e: return json.dumps({ "error": f"知识库查询失败: {e}", "question": question, }, ensure_ascii=False) # ============================================================ # 需风量数据查询(通防管控平台 MCP) # ============================================================ async def get_needq_all_data() -> str: """从通防管控平台获取全部用风地点的需风量数据。 调用远程 MCP 工具 get_needq_all_data,返回管控平台上所有用风地点的 计划需风量、实际风量、偏差等数据。用于回答用户"查看管控平台的需风量情况" "通防管控平台上各地点需风量是多少"等问题。 Returns: JSON 格式的全部需风量数据,包含各用风地点的名称、类型、计划需风量、 实际风量、偏差等字段。 """ return await _call_mcp_tool("get_needq_all_data", {}) # ============================================================ # 监测历史数据查询(通防管控平台 MCP) # ============================================================ async def list_ventanaly_monitor_data_days( strtype: str = "", gdeviceids: str = "", ttime_begin: str = "", ttime_end: str = "", device_num: str = "", skip: int = 8, page_no: int = 1, page_size: int = 100, # column: str = "", ) -> str: """查询监测设备历史时序数据。 根据设备类型、设备ID、时间范围等条件,查询监测设备的历史时序数据。 适用于查询风速、风量、瓦斯、温度等传感器在一段时间内的历史趋势。 Args: strtype: 设备类型(必填),对应后端 strtype,如 "fanmain_stem_wp_2" gdeviceids: 设备ID(必填),如 "11111004";多个设备按后端要求格式拼接 ttime_begin: 开始时间(必填),格式 "yyyy-MM-dd HH:mm:ss" ttime_end: 结束时间(必填),格式 "yyyy-MM-dd HH:mm:ss" device_num: 设备编号(必填),对应后端 deviceNum,如 "Fan1" skip: 采样间隔/跳点参数(必填),默认8,1=5秒//2=10秒//3=30秒//4=1分钟//5=5分钟//6=10分钟//7=30分钟8=1小时,如果半天内默认按10分钟查询,如果需要跨天则默认按1小时查询,其他情况请按合适的采样间隔获取,考虑数据库的压力。 page_no: 页码,默认 1 page_size: 每页条数,默认 200 column: 排序字段,可选 Returns: JSON格式的历史监测时序数据。 """ return await _call_mcp_tool("list_ventanaly_monitor_data_days", { "strtype": strtype, "gdeviceids": gdeviceids, "ttime_begin": ttime_begin, "ttime_end": ttime_end, "device_num": device_num, "skip": skip, "page_no": page_no, "page_size": page_size, # "column": column, })