From 4989a0898f9dbba51d073e61a479021d96f5b5e6 Mon Sep 17 00:00:00 2001 From: dimitrievgs Date: Sun, 14 Jun 2026 21:06:09 +0300 Subject: [PATCH] Better display of exception, better parsing of env parameter, better collection of req and res for tool --- app/mcp_tools.py | 65 +++++++++++++++++++++++++++++++++++++++++------- app/workflows.py | 56 ++++++++++++++++++++++++++++------------- requirements.txt | 2 ++ 3 files changed, 97 insertions(+), 26 deletions(-) diff --git a/app/mcp_tools.py b/app/mcp_tools.py index 93772f7..3b1a950 100644 --- a/app/mcp_tools.py +++ b/app/mcp_tools.py @@ -1,5 +1,5 @@ from pydantic import BaseModel, Field, create_model -from langchain_core.tools import StructuredTool +from langchain_core.tools import StructuredTool, ToolException # ----------------- MCP MANAGER ----------------- try: @@ -54,17 +54,55 @@ def create_mcp_tool(server_config, tool_name, tool_desc, json_schema, full_env, await session.initialize() result = await session.call_tool(tool_name, arguments=kwargs) text_outputs = [] + + # Проверяем, не вернул ли сам MCP-пакет флаг ошибки (isError) + # Если да, мы тоже должны это воспринимать как падение + is_mcp_error = getattr(result, 'isError', False) + for c in result.content: if hasattr(c, 'text'): text_outputs.append(c.text) else: text_outputs.append(str(c)) - return "\n".join(text_outputs) + + output_text = "\n".join(text_outputs) + + # Если сервер MCP явно сказал об ошибке, кидаем Exception + if is_mcp_error: + raise Exception(output_text) + + return output_text - raw_output = asyncio.run(_run()) + # Оборачиваем попытку запуска в try..except + try: + raw_output = asyncio.run(_run()) + except Exception as e: + def extract_root_errors(exc): + # Если ошибка содержит вложенные ошибки (ExceptionGroup) + if hasattr(exc, 'exceptions'): + msgs = [] + for child_exc in exc.exceptions: + msgs.append(extract_root_errors(child_exc)) + # Объединяем сообщения, убирая пустые + return " | ".join(filter(bool, msgs)) + # Если это базовая ошибка (например, McpError), возвращаем ее текст + return str(exc) + + clean_error_msg = extract_root_errors(e) + # 1. Поймали обрыв связи (сервер не запущен, отвалился stdio и т.д.) + error_msg = f"Отсутствует связь с сервером или ошибка выполнения: {clean_error_msg}" + + # 2. Фиксируем kwargs (запрос), чтобы фронтенд в блоке Response/Request показал JSON! + if debug_callback: + debug_callback(kwargs, error_msg) + + # 3. ВАЖНО: выбрасываем ошибку дальше. Бэкенд LangChain ее перехватит + # и отправит на фронтенд SSE-событие с `is_error: true` -> появится КРАСНЫЙ КРЕСТ. + raise ToolException(error_msg) + + # Выполняется только при успехе (зеленый чекмарк) if debug_callback: - # debug_callback(f"MCP: {tool_name}\n" + input_json, raw_output) debug_callback(kwargs, raw_output) return f"Инструмент '{tool_name}' вернул следующий результат:\n\n{raw_output}\n" @@ -73,7 +111,8 @@ def create_mcp_tool(server_config, tool_name, tool_desc, json_schema, full_env, func=mcp_tool_runner, name=tool_name.replace('-','_'), # Langchain не любит тире в именах description=tool_desc or f"MCP Tool {tool_name}", - args_schema=args_schema + args_schema=args_schema, + handle_tool_error=True # Разрешает агенту "выжить" после ошибки тула и сгенерировать ответ ) def fetch_mcp_tools(server_config, debug_callback=None): @@ -98,10 +137,18 @@ def fetch_mcp_tools(server_config, debug_callback=None): env_dict = {} env_str = server_config.get("envString", "") if env_str: - for pair in env_str.split(','): - if '=' in pair: - k, v = pair.split('=', 1) - env_dict[k.strip()] = v.strip() + import re + # Регулярка ищет: КЛЮЧ = "ЗНАЧЕНИЕ" | 'ЗНАЧЕНИЕ' | ЗНАЧЕНИЕ_ДО_ЗАПЯТОЙ + pattern = r'([^,= \t]+)\s*=\s*(?:"([^"]*)"|\'([^\']*)\'|([^,]*))' + + for m in re.finditer(pattern, env_str): + k = m.group(1) + # Берем то совпадение, которое сработало (2 - двойные кавычки, 3 - одинарные, 4 - без кавычек) + v = m.group(2) if m.group(2) is not None else \ + m.group(3) if m.group(3) is not None else \ + m.group(4) if m.group(4) is not None else "" + + env_dict[k.strip()] = v.strip() full_env = {**os.environ.copy(), **env_dict} server_params = StdioServerParameters(command=command, args=args, env=full_env) diff --git a/app/workflows.py b/app/workflows.py index 51766c1..5b74fa4 100644 --- a/app/workflows.py +++ b/app/workflows.py @@ -100,7 +100,7 @@ def run_agent_streaming(graph_id: str, elif role == "assistant": messages_for_llm.append(AIMessage(content=msg.get("content", ""))) - # Локальное хранилище для отладки текущего запуска + # Локальное хранилище для отладки текущего запуска (на случай кастомных тулов) current_run_debug = [] def capture_debug(req, res): @@ -116,30 +116,45 @@ def run_agent_streaming(graph_id: str, debug_callback=capture_debug ) - # 4. Выполняем цикл с потоковой передачей событий + # Буфер для вызовов тулов от LLM step_counter = 0 - tool_mapping = {} # Хранит связку: ID вызова -> Номер шага + tool_mapping = {} # Хранит связку: tc_id -> Номер шага + active_tool_calls = {} # Хранит прогресс для конкретных инструментов { tc_id: {"name": ..., "args": ...} } + current_idx_map = {} # Временный маппинг: index (из чанков текущего ответа) -> глобальный tc_id - # LangGraph Stream Mode 'messages': отдает чанки токенов и вызовы функций + # 4. Выполняем цикл с потоковой передачей событий for chunk, metadata in agent_executor.stream({"messages": messages_for_llm}, stream_mode="messages"): if isinstance(chunk, AIMessageChunk): # Агент запросил вызов инструмента (Формирование JSON аргументов) if hasattr(chunk, "tool_call_chunks") and chunk.tool_call_chunks: for tc in chunk.tool_call_chunks: - if "name" in tc and tc["name"]: + idx = tc.get("index", 0) + + # При поступлении самого первого чанка для инструмента, в нем есть УНИКАЛЬНЫЙ 'id' + if "id" in tc and tc["id"]: + tc_id = tc["id"] + # Запоминаем, что индекс idx в текущем цикле принадлежит уникальному tc_id + current_idx_map[idx] = tc_id + step_counter += 1 - tc_id = tc.get("id", f"unknown_{step_counter}") - # Сохраняем шаг для этого конкретного вызова tool_mapping[tc_id] = step_counter + active_tool_calls[tc_id] = {"name": tc.get("name", "unknown_tool"), "args": ""} yield { "type": "tool_start", - "name": tc["name"], + "name": active_tool_calls[tc_id]["name"], "step": step_counter, "node_id": assistant_node_id } + # Считываем сырой JSON кусочек за кусочком по мере их поступления + if "args" in tc and tc["args"]: + # Достаем сохраненный уникальный id из маппинга для текущего сообщения + tc_id = current_idx_map.get(idx) + if tc_id and tc_id in active_tool_calls: + active_tool_calls[tc_id]["args"] += tc["args"] + # Обычный стриминг текста финального ответа if chunk.content: yield { @@ -150,22 +165,29 @@ def run_agent_streaming(graph_id: str, # Тул отработал и принес ответ elif hasattr(chunk, "type") and chunk.type == "tool": - req = "" res = str(chunk.content) + tc_id = getattr(chunk, "tool_call_id", "") - # Достаем JSON отправленных параметров из последнего debug_callback - if current_run_debug: - dbg = current_run_debug.pop(0) # Берем первый из очереди - req = dbg.get("req", "") + # 1. Забираем готовую строку (JSON) из нашего прямого перехвата LLM по уникальному tc_id! + req = active_tool_calls.get(tc_id, {}).get("args", "") + # 2. Резервный механизм на старых колбеках (и форматирование в json если это dict) + if not req and current_run_debug: + dbg = current_run_debug.pop(0) + rq = dbg.get("req", "") + if isinstance(rq, dict): + import json + req = json.dumps(rq, indent=2, ensure_ascii=False) + else: + req = str(rq) + # Ловушка для определения ошибок is_error = False if getattr(chunk, "status", "") == "error" or "ОШИБКА ИСПОЛНЕНИЯ" in res or "SCRIPT RUNTIME ERROR" in res: is_error = True - - # Достаем ИМЕННО ТОТ ШАГ, на котором начинался вызов этого инструмента - tc_id = getattr(chunk, "tool_call_id", "") - actual_step = tool_mapping.get(tc_id, step_counter) # Fallback, если id нет + + # Ищем шаг тула + actual_step = tool_mapping.get(tc_id, step_counter) yield { "type": "tool_end", diff --git a/requirements.txt b/requirements.txt index 3907cc5..5048f9c 100644 --- a/requirements.txt +++ b/requirements.txt @@ -6,6 +6,8 @@ langchain_core==0.3.76 langchain_google_genai==2.1.12 +mcp==1.27.2 + flask==3.1.2 flask_cors==6.0.1 flask_socketio