# -*- 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: 查询监测历史时序数据 - get_device_kind_dict: 查询设备大类小类全量字典 - get_device_list_by_kind: 根据设备类型查询设备列表 - query_device_realtime_data: 查询设备实时监测数据 """ import json import contextvars from fastmcp import Client from tools.mcp_logger import log_mcp_call # ============================================================ # 当前用户上下文(供工具函数获取调用者身份) # ============================================================ _current_user: contextvars.ContextVar[str] = contextvars.ContextVar( 'current_user', default='admin' ) # ============================================================ # 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", {}) async def get_device_kind_dict() -> str: """查询设备大类(deviceKind)和小类(strType)的全量字典。 返回设备类型编码与中文名称的完整映射表,包含 deviceKind(设备大类) 和 strType(设备小类)两个维度的编码-名称对照。 无参数。 适用场景: - 用户询问"有哪些设备类型""设备分类有哪些""设备大类小类有哪些" - 需要将设备类型编码翻译为中文名称 - 查询某种设备类型对应的 strType 编码以便后续调用历史数据工具 - 用户想知道系统中支持监控哪些类型的设备 Returns: JSON 格式的设备类型字典,包含 deviceKind 和 strType 映射。 """ return await _call_mcp_tool("get_device_kind_dict", {}) async def get_device_list_by_kind(device_kind: str) -> str: """根据设备类型(deviceKind)查询该类型下的全量设备列表。 返回指定设备大类下的所有设备,包含设备ID、设备名称、安装位置、分站名称。 适用场景: - 用户询问"列出所有风速传感器""有哪些甲烷传感器"等按类型筛选设备 - 需要获取某类设备的完整清单以便进一步查询实时/历史数据 - 结合 get_device_kind_dict 先获取类型编码,再按类型查设备列表 Args: device_kind: 设备大类编码(必填),如 "modelsensor_speed"表示风速传感器。 可先调用 get_device_kind_dict 获取所有可用的 deviceKind 编码。 Returns: JSON 格式的设备列表,包含 device_id、device_name、install_pos、station_name。 """ return await _call_mcp_tool("get_device_list_by_kind", {"device_kind": device_kind}) async def query_device_realtime_data(device_id: str) -> str: """查询指定设备的实时监测数据。 根据设备ID获取该设备当前最新的监测读数、运行状态和报警信息。 与 query_device_data_by_id 不同,本工具专注于单设备的实时快照数据, 返回结构更精简,延迟更低。 适用场景: - 用户询问"设备XXX的实时数据""传感器XXX当前读数是多少" - 已知设备ID,需要快速获取其最新监测值 - 配合 get_device_list_by_kind 先查出设备ID列表,再逐个查询实时数据 Args: device_id: 设备ID(必填),可从 query_devices_by_tunnel / get_device_list_by_kind 返回结果中获取。 Returns: JSON 格式的设备实时监测数据,包含当前读数、单位、采集时间、在线状态、报警信息。 """ return await _call_mcp_tool("query_device_realtime_data", {"device_id": device_id}) # ============================================================ # 监测历史数据查询(通防管控平台 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, }) # ============================================================ # 用户偏好记忆工具 # ============================================================ def save_user_preference(content: str, keywords: str = "", category: str = "通用") -> str: """保存用户偏好/习惯到个人记忆库。 当用户在对话中明确要求"记住""保存为习惯/偏好""以后都用这个"时调用。 下次该用户对话时,系统会自动注入已保存的偏好作为上下文。 Args: content: 偏好内容,描述具体的习惯或个性化要求。例如"习惯使用 m³/s 而非 m³/min"、"15216工作面默认采高3.5m" keywords: 触发关键词,多个用逗号分隔。例如"风速,单位,风量"。留空则自动匹配。 category: 偏好分类,默认"通用"。可选值:需风量计算、数据解读、规程查询、通用 Returns: 保存结果,含记录ID供后续删除用。 """ from db.chat_store import save_user_preference as _db_save user_name = _current_user.get() pref_id = _db_save(user_name, content, keywords, category) return json.dumps({ "success": True, "id": pref_id, "message": f"已保存偏好 (id={pref_id}):{content}", }, ensure_ascii=False) def list_user_preferences() -> str: """查看当前用户已保存的所有偏好/习惯。 列出该用户所有偏好记录,包含ID、分类、内容、关键词、保存时间。 Returns: JSON格式的偏好列表。 """ from db.chat_store import get_user_preferences as _db_list user_name = _current_user.get() prefs = _db_list(user_name) if not prefs: return json.dumps({"preferences": [], "message": "暂无保存的偏好"}, ensure_ascii=False) return json.dumps({ "preferences": [ {"id": p["id"], "category": p["category"], "content": p["content"], "keywords": p["keywords"], "created_at": p["created_at"]} for p in prefs ], }, ensure_ascii=False) def delete_user_preference(preference_id: int) -> str: """删除一条用户偏好记录。 Args: preference_id: 要删除的偏好记录ID(可从 list_user_preferences 获取) Returns: 删除结果。 """ from db.chat_store import delete_user_preference as _db_delete user_name = _current_user.get() ok = _db_delete(preference_id, user_name) if ok: return json.dumps({"success": True, "message": f"已删除偏好 (id={preference_id})"}, ensure_ascii=False) return json.dumps({"success": False, "message": f"未找到偏好记录 id={preference_id} 或无权操作"}, ensure_ascii=False) # ============================================================ # 计划审批工具(Human-in-the-Loop) # ============================================================ def request_plan_approval(plan_summary: str) -> str: """提交执行计划等待人工审批。在制定好完整计划后调用此工具。 仅在计划模式(plan mode)下由 Agent 主动调用。 调用后会暂停执行,等待用户在前端审批(批准/拒绝)。 审批通过后自动继续执行计划。 Args: plan_summary: 执行计划的简要描述,需包含: - 计划分几步,每步做什么 - 每步预期调用哪些工具 - 预期的输出结果 Returns: "计划已批准,开始执行。" 或 "计划被拒绝。" """ from langgraph.types import interrupt result = interrupt({ "type": "plan_approval", "plan": plan_summary, "message": "智能体已制定执行计划,等待您的审批...", }) if isinstance(result, dict) and result.get("action") == "approve": return "计划已批准,开始执行。" else: return "计划被拒绝。"