""" Этот файл содержит API endpoints, созданные с использованием Flask. Он обеспечивает взаимодействие с агентом через HTTP запросы, включая обработку сообщений, управление графами и получение истории. """ import json from flask import Flask, request, jsonify, Response, stream_with_context from flask_cors import CORS from flask_socketio import SocketIO from workflows import graph_history_manager, run_agent_streaming from llm_client import DEFAULT_TEMPERATURE, MODELS, get_llm # Добавляем импорт списка моделей import base64 from mcp_tools import get_raw_mcp_tools_list import threading import uuid as uuid_lib api = Flask(__name__) CORS( api ) # Разрешаем CORS для всех доменов (в production нужно настроить более строго) # ----------------------------------------------- WebSocket Initialization --------------------------------------------------------------- socketio = SocketIO( api, cors_allowed_origins="*") # В production настроить CORS строже # , async_mode='eventlet' для компиляции в exe graph_history_manager.set_socketio(socketio) # ----------------------------------------------- WebSocket Events --------------------------------------------------------------- @socketio.on('connect') def handle_connect(): """Обработчик подключения клиента к WebSocket.""" print(f"Client connected: {request.sid}") @socketio.on('disconnect') def handle_disconnect(): """Обработчик отключения клиента от WebSocket.""" print(f"Client disconnected: {request.sid}") @socketio.on('client_proactive_ping') def handle_client_proactive_ping(data=None): """Обработчик проактивного пинга от клиента.""" print(f"Received proactive ping from client: {request.sid}. Sending pong.") # Возвращаем словарь, который будет преобразован в JSON и отправлен как ответ на callback return {'status': 'pong'} # ----------------------------------------------- API Endpoints --------------------------------------------------------------- @api.route('/api/graphs', methods=['GET']) def get_graphs(): """API endpoint для получения списка графов.""" graphs = graph_history_manager.get_all_graphs_summary() print("get_graphs response:", graphs) # Добавлено логирование return jsonify(graphs) @api.route('/api/graphs/new', methods=['POST']) def create_new_graph_endpoint(): """API endpoint для создания нового пустого графа.""" try: new_graph_id = graph_history_manager.create_new_graph() return jsonify({"graph_id": new_graph_id}), 201 except Exception as e: print(f"Ошибка при создании нового графа через API: {e}") return jsonify({"error": f"Не удалось создать новый граф: {str(e)}"}, 500) @api.route('/api/graphs/', methods=['GET']) def get_graph_by_id(graph_id): """API endpoint для получения данных конкретного графа по ID.""" graph_data = graph_history_manager.get_graph(graph_id) if graph_data: return jsonify(graph_data) else: return jsonify({"error": f"Граф с ID {graph_id} не найден."}, 404) @api.route('/api/graphs/', methods=['DELETE']) def delete_graph(graph_id): """API endpoint для удаления графа.""" if not graph_id: return jsonify({"error": "Не указан ID графа для удаления."}, 400) if graph_history_manager.delete_graph(graph_id): return jsonify({"message": f"Граф {graph_id} успешно удален."}) else: return jsonify({"error": f"Не удалось удалить граф {graph_id}."}, 500) @api.route('/api/graphs', methods=['DELETE']) def delete_all_graphs(): """API endpoint для удаления всех графов.""" if graph_history_manager.delete_all_graphs(): return jsonify({"message": "Все графы успешно удалены."}) else: return jsonify({"error": "Не удалось удалить все графы."}, 500) @api.route('/api/graphs//rename', methods=['PATCH']) def rename_graph(graph_id): """API endpoint для переименования графа.""" if not graph_id: return jsonify({"error": "Не указан ID графа для переименования."}, 400) data = request.get_json() new_title = data.get("title") if not new_title or not new_title.strip(): return jsonify({"error": "Название не может быть пустым."}, 400) if graph_history_manager.rename_graph(graph_id, new_title.strip()): return jsonify({"message": f"Граф {graph_id} успешно переименован."}) else: return jsonify({"error": f"Не удалось переименовать граф {graph_id}."}, 500) @api.route('/api/graphs//nodes/', methods=['DELETE']) def delete_graph_node(graph_id, node_id): """API endpoint для удаления узла графа.""" if not graph_id or not node_id: return jsonify({"error": "Не указаны ID графа или узла для удаления."}, 400) try: if graph_history_manager.delete_node(graph_id, node_id): return jsonify({ "message": f"Узел {node_id} успешно удален из графа {graph_id}." }) else: return jsonify( { "error": f"Не удалось удалить узел {node_id} из графа {graph_id}." }, 500) except Exception as e: print(f"Ошибка при удалении узла {node_id} из графа {graph_id}: {e}") return jsonify( {"error": f"Ошибка сервера при удалении узла: {str(e)}"}, 500) @api.route('/api/graphs//settings', methods=['PATCH']) def update_graph_settings_endpoint(graph_id): """API endpoint для обновления настроек графа (например, пользовательского системного промпта).""" data = request.get_json() custom_system_prompt = data.get( "custom_system_prompt") # Может быть None для сброса if not graph_id: return jsonify( {"error": "Не указан ID графа для обновления настроек."}, 400) try: if graph_history_manager.update_graph_settings(graph_id, custom_system_prompt): return jsonify( {"message": f"Настройки графа {graph_id} успешно обновлены."}) else: return jsonify( {"error": f"Не удалось обновить настройки графа {graph_id}."}, 500) except Exception as e: print(f"Ошибка при обновлении настроек графа {graph_id}: {e}") return jsonify({"error": f"Ошибка сервера: {str(e)}"}, 500) @api.route('/api/messages//', methods=['GET']) def get_messages_from_root_to_node(graph_id, node_id): """API endpoint для получения сообщений от корня до выбранной ноды.""" graph_data = graph_history_manager.get_graph(graph_id) if not graph_data: return jsonify({"error": "Граф не найден."}, 404) target_node_id = node_id if node_id == 'last': target_node_id = graph_data.get('current_node_id') if not target_node_id: # Если запрошен 'last', но активного узла нет, возвращаем пустой список return jsonify([]) messages = graph_history_manager.get_messages_from_root_to_node( graph_id, graph_data, target_node_id) return jsonify(messages) @api.route('/api/models', methods=['GET']) def get_available_models(): """API endpoint для получения списка доступных моделей.""" return jsonify(list(MODELS.keys())) @api.route('/api/chat/send', methods=['POST']) def send_user_message(): """API endpoint для создания узла пользователя с вложениями.""" data = request.get_json() message = data.get("message", "") graph_id = data.get("graph_id") parent_node_id = data.get("parent_node_id") attachments = data.get("attachments", []) cache_folder = data.get("cache_folder") # Блокируем только если нет и текста, и вложений if not message and not attachments: return jsonify({"error": "Сообщение или вложение не могут быть пустыми."}), 400 try: result = graph_history_manager.create_user_node( message, graph_id, parent_node_id, attachments, cache_folder ) graph_history_manager.title_generator.add_node_to_queue_direct( result.get("graph_id"), result.get("node_id")) return jsonify(result) except Exception as e: print(f"Ошибка при создании узла пользователя: {e}") return jsonify({"error": str(e)}, 500) @api.route('/api/chat/stream', methods=['POST']) def chat_stream(): """API endpoint для стриминга ответа LLM и работы ReAct-агента.""" data = request.get_json() graph_id = data.get("graph_id") user_node_id = data.get("user_node_id") system_prompt = data.get("system_prompt") model = data.get("model") cache_folder = data.get("cache_folder") temperature = data.get("temperature", DEFAULT_TEMPERATURE) max_tokens = data.get("max_tokens") agency_mode = data.get("agency_mode", False) # Необязательный параметр для существующего узла ассистента при регенерации existing_assistant_node_id = data.get("existing_assistant_node_id") if not graph_id or not user_node_id: return jsonify({"error": "Не указаны graph_id или user_node_id."}, 400) # Подтягиваем глобальные настройки Obsidian (синхронизируются через /settings/sync) global current_obsidian_settings def generate(): assistant_node_id = None try: if existing_assistant_node_id: assistant_node_id = existing_assistant_node_id graph_history_manager.mark_node_as_placeholder_and_clear_content( graph_id, assistant_node_id, user_node_id) yield f"data: {json.dumps({'type': 'node_created', 'node_id': assistant_node_id, 'reused': True})}\n\n" else: # Создаём узел-заглушку для ответа LLM assistant_node_id = graph_history_manager.create_assistant_placeholder_node( graph_id, user_node_id) # Отправляем ID нового узла клиенту yield f"data: {json.dumps({'type': 'node_created', 'node_id': assistant_node_id, 'reused': False})}\n\n" # Запускаем стриминг ответа от LLM (используем обновленный агент) accumulated_content = "" # Вспомогательная функция для безопасной передачи JSON в HTML-атрибутах def safe_b64(text): if not text: return "e30=" # пустой json '{}' return base64.b64encode(str(text).encode('utf-8')).decode('utf-8') try: for chunk_data in run_agent_streaming(graph_id, user_node_id, assistant_node_id, system_prompt, model, cache_folder, temperature, max_tokens, agency_mode, current_obsidian_settings): # Прокидываем любые успешные этапы: куски текста или статусы работы тулов if chunk_data.get("type") in ["chunk", "tool_start", "tool_end"]: if chunk_data.get("type") == "chunk": accumulated_content += chunk_data.get("content", "") elif chunk_data.get("type") == "tool_start": # Встраиваем стартовый маркер в память базы данных start_marker = f'\n\n
\n\n' accumulated_content += start_marker elif chunk_data.get("type") == "tool_end": # Находим стартовый маркер и меняем его на финальный с данными search_marker = f'
' req_b64 = safe_b64(chunk_data.get("request")) res_b64 = safe_b64(chunk_data.get("response")) status = "error" if chunk_data.get("is_error") else "success" end_marker = f'
' if search_marker in accumulated_content: accumulated_content = accumulated_content.replace(search_marker, end_marker) else: accumulated_content += f'\n\n{end_marker}\n\n' yield f"data: {json.dumps(chunk_data)}\n\n" # Обработка ошибки во время стриминга (крашим узел на фронте) elif chunk_data.get("type") == "error": if assistant_node_id: graph_history_manager.mark_node_as_error( graph_id, assistant_node_id, chunk_data.get("error", "Неизвестная ошибка при стриминге")) if graph_history_manager.socketio: graph_history_manager.socketio.emit( 'node_type_changed', { 'graph_id': graph_id, 'node_id': assistant_node_id, 'node_type': 'error' }) yield f"data: {json.dumps(chunk_data)}\n\n" return except GeneratorExit: # Обработка прерывания генератора (если клиент нажал Stop на фронте или разорвал коннект) print(f"⚠️ Клиент отменил стрим для узла {assistant_node_id}") # После завершения стриминга (успешного) обновляем узел полным контентом в базе SQLite if assistant_node_id: graph_history_manager.update_assistant_node_content( graph_id, assistant_node_id, accumulated_content) # Инициируем фоновую генерацию заголовка graph_history_manager.title_generator.add_node_to_queue_direct( graph_id, assistant_node_id) yield f"data: {json.dumps({'type': 'done', 'node_id': assistant_node_id})}\n\n" except Exception as e: # Отлов непредсказуемых внешних ошибок в самом генераторе API print(f"Ошибка при стриминге ответа: {e}") error_message = str(e) if assistant_node_id: graph_history_manager.mark_node_as_error( graph_id, assistant_node_id, f"Ошибка при стриминге: {error_message}") if graph_history_manager.socketio: graph_history_manager.socketio.emit( 'node_type_changed', { 'graph_id': graph_id, 'node_id': assistant_node_id, 'node_type': 'error' }) error_data = { 'type': 'error', 'error': f"Ошибка при стриминге: {error_message}" } yield f"data: {json.dumps(error_data)}\n\n" return Response(stream_with_context(generate()), content_type='text/event-stream') @api.route('/api/chat/regenerate', methods=['POST']) def regenerate_message(): """API endpoint для регенерации сообщения. При ошибке LLM-узла, переиспользует его. В противном случае, создает новый узел. """ data = request.get_json() graph_id = data.get("graph_id") node_id = data.get( "node_id") # Это ID узла, который нужно регенерировать (llm-ответ) model = data.get("model") system_prompt = data.get("system_prompt") cache_folder = data.get("cache_folder") temperature = data.get("temperature", DEFAULT_TEMPERATURE) max_tokens = data.get("max_tokens") if not graph_id or not node_id: return jsonify({"error": "Не указаны graph_id или node_id."}, 400) try: # Определяем родительский узел для LLM-ответа. Это должен быть user-узел. parent_node_id = graph_history_manager.get_parent_node_id( graph_id, node_id) if not parent_node_id: return jsonify( { "error": "Не удалось найти родительский узел для регенерации." }, 404) # НОВОЕ: Проверяем тип узла node_type = graph_history_manager.get_node_type(graph_id, node_id) assistant_node_id_to_use = None if node_type == 'error': # Если узел был ошибочным, переиспользуем его # Очищаем его содержимое и помечаем как заглушку сразу же assistant_node_id_to_use = node_id # Используем тот же ID else: # Если узел не был ошибочным, создаем новый узел LLM-ответа assistant_node_id_to_use = graph_history_manager.create_assistant_placeholder_node( graph_id, parent_node_id) return jsonify({ "existing_assistant_node_id": assistant_node_id_to_use, "user_node_id": parent_node_id, "graph_id": graph_id, "model": model, "system_prompt": system_prompt, "cache_folder": cache_folder, "temperature": temperature, "max_tokens": max_tokens }) except Exception as e: print(f"Ошибка при регенерации сообщения: {e}") return jsonify({"error": str(e)}, 500) @api.route('/api/mcp/fetch', methods=['POST']) def fetch_mcp_tools_route(): """API endpoint для получения списка инструментов MCP сервера (без создания LangChain тулов)""" server_config = request.get_json() if not server_config: return jsonify({"error": "Пустой конфиг сервера"}), 400 try: vault_cwd = current_obsidian_settings.get("vaultAbsolutePath") tools_list = get_raw_mcp_tools_list(server_config, vault_cwd=vault_cwd) return jsonify(tools_list) except Exception as e: return jsonify({"error": str(e)}), 500 from voice_service import VoiceService, VOICE_COMMANDS_RESPONSE_TO_STORE # 1. Создаем экземпляр (обязательно укажите путь к вашей модели Vosk) voice_inst = VoiceService() current_obsidian_settings = {} obsidian_settings_fetched = False @api.route('/api/voice/status', methods=['GET']) def get_voice_status(): # 2. Вызываем метод у экземпляра voice_inst, а не у модуля pure_gm_text, session_id = voice_inst.get_pure_gm_text() return jsonify({ "transcription": voice_inst.get_context_by_limit(current_obsidian_settings.get('voiceDisplayLimit', 0), order='reversed'), "pure_gm_text": pure_gm_text, "session_id": session_id, "response1": list(voice_inst.responses["regular"]), "response2": list(voice_inst.responses["commands"]), "metrics": voice_inst.metrics_data, "is_running": voice_inst.is_running, "obsidian_settings_fetched": obsidian_settings_fetched }) @api.route('/api/voice/control', methods=['POST']) def control_voice(): action = request.json.get("action") if action == "start": voice_inst.start_session(False) elif action == "continue": voice_inst.start_session(True) elif action == "stop": voice_inst.stop_session() return jsonify({"status": "ok"}) @api.route('/api/settings/sync', methods=['POST']) def sync_settings(): global current_obsidian_settings, obsidian_settings_fetched obsidian_settings_fetched = True data = request.json if data: current_obsidian_settings.update(data) # Передаем настройки в voice_service if data.get('absoluteLogsPath'): # Если фронт прислал вычисленный абсолютный путь voice_inst.update_logs_config( abs_path=data.get('absoluteLogsPath') ) # Инициализируем модель Vosk, если передан путь vosk_path = data.get('voskModelPath') if vosk_path: # Запускаем в фоновом потоке, чтобы не подвесить фронтенд Obsidian (загрузка занимает время) import threading threading.Thread(target=voice_inst.load_model, args=(vosk_path,), daemon=True).start() # TODO: Скорее всего, здесь нужно какое-то событие эмитировать об изменении настроек, и в другом месте ловить graph_history_manager.title_generator.llm = get_llm(current_obsidian_settings['summarizationModel']) graph_history_manager.title_generator.voice_llm = get_llm(current_obsidian_settings['voiceCommandModel']) print(f"✅⚙️ Settings synced: {current_obsidian_settings}") return jsonify({"status": "synced"}) # --------------------------------------------------------------------------- # Vault Query Bridge # Позволяет Python-скриптам читать/писать файлы vault через Obsidian. # Схема: POST /api/vault/query → socketio emit → Obsidian отвечает → возврат. # --------------------------------------------------------------------------- import threading import uuid as uuid_lib # Хранилище ожидающих ответа запросов. # Ключ: request_id (str UUID) # Значение: (threading.Event, dict результата) _vault_query_pending: dict[str, tuple[threading.Event, dict]] = {} @api.route('/api/vault/query', methods=['POST']) def vault_query_endpoint(): """ Проксирует запрос к Obsidian vault через WebSocket. Принимает JSON: { "type": "get_frontmatter" | "get_all_factions" | ..., "payload": { ... параметры операции ... } } Поддерживаемые типы (обрабатываются на стороне Obsidian VaultQueryHandler): get_frontmatter — YAML + тело одного файла get_all_factions — все файлы в папке (с фильтром по тегу) get_recent_files — последние N файлов (с фильтром по тегам) resolve_wikilink — wikilink → путь к файлу write_frontmatter — перезаписать YAML + тело файла create_file — создать или перезаписать файл Блокируется до ответа Obsidian или до истечения timeout (10с). При timeout возвращает 504. """ data = request.get_json(silent=True) if not data: return jsonify({"error": "Пустое тело запроса"}), 400 query_type = data.get("type") payload = data.get("payload", {}) if not query_type: return jsonify({"error": "Поле 'type' обязательно"}), 400 # Генерируем уникальный ID для сопоставления запрос↔ответ request_id = str(uuid_lib.uuid4()) # Регистрируем ожидание ответа event = threading.Event() result_holder: dict = {} _vault_query_pending[request_id] = (event, result_holder) try: # Отправляем запрос в Obsidian через WebSocket socketio.emit('vault_query_request', { "request_id": request_id, "type": query_type, "payload": payload }) # Блокируемся — ждём пока handle_vault_query_response вызовет event.set() answered = event.wait(timeout=10.0) finally: # Убираем из pending независимо от результата _vault_query_pending.pop(request_id, None) if not answered: return jsonify({ "error": ( "Timeout: Obsidian не ответил за 10 секунд. " "Проверьте что плагин LLM Agent запущен и подключён к WebSocket." ) }), 504 if "error" in result_holder: return jsonify({"error": result_holder["error"]}), 500 return jsonify(result_holder.get("data")) @socketio.on('vault_query_response') def handle_vault_query_response(data): """ Получает ответ от Obsidian на vault-запрос. Будит заблокированный поток в vault_query_endpoint(). data ожидается вида: { "request_id": "...", "result": <любые данные> } или { "request_id": "...", "error": "описание ошибки" } """ if not isinstance(data, dict): print(f"⚠️ vault_query_response: неожиданный формат: {type(data)}") return request_id = data.get("request_id") if not request_id: print("⚠️ vault_query_response: отсутствует request_id") return pending = _vault_query_pending.get(request_id) if not pending: # Дубль после таймаута или успешной обработки — это нормально # если исправлен WebSocketService. До исправления — признак бага. print(f"⚠️ дубль vault_query_response [{request_id[:8]}] — игнорируем") return event, result_holder = pending if "error" in data: result_holder["error"] = data["error"] else: result_holder["data"] = data.get("result") print(f"✅ vault_query_response [{request_id[:8]}]: тип={type(result_holder['data']).__name__}") event.set() # pending будет удалён в finally блока vault_query_endpoint