| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335 |
- # -*- 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,
- })
|