From 196201c681bd86fd91dbc9debcc762caeb95f938 Mon Sep 17 00:00:00 2001 From: dimitrievgs Date: Sun, 12 Oct 2025 21:40:50 +0300 Subject: [PATCH] More correct working with error nodes + more correct regeneration --- .gitignore | 2 +- app/api.py | 99 +++++++++++++++++++++++++++--------- app/graph_history_manager.py | 56 ++++++++++++++++++++ app/llm_client.py | 10 ++++ app/nodes.py | 2 +- app/title_generator.py | 66 ++++++++++++++++++++---- 6 files changed, 199 insertions(+), 36 deletions(-) diff --git a/.gitignore b/.gitignore index b2e4d9a..e94805f 100644 --- a/.gitignore +++ b/.gitignore @@ -6,4 +6,4 @@ node_modules/ __pycache__ *.db -*.db-history \ No newline at end of file +*.db-journal \ No newline at end of file diff --git a/app/api.py b/app/api.py index 1a0e538..d014059 100644 --- a/app/api.py +++ b/app/api.py @@ -35,6 +35,13 @@ 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']) @@ -118,6 +125,9 @@ def send_user_message(): try: result = graph_history_manager.create_user_node(message, graph_id, parent_node_id) + + 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}") @@ -132,17 +142,27 @@ def chat_stream(): user_node_id = data.get("user_node_id") system_prompt = data.get("system_prompt") model = data.get("model") + # Необязательный параметр для существующего узла ассистента при регенерации + 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) def generate(): + assistant_node_id = None try: - # Создаём узел-заглушку для ответа 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})}\n\n" + 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 = "" @@ -151,20 +171,39 @@ def chat_stream(): accumulated_content += chunk_data.get("content", "") 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 + return # Прекращаем генерацию при ошибке # После завершения стриминга обновляем узел полным контентом - graph_history_manager.update_assistant_node_content(graph_id, assistant_node_id, accumulated_content) + 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) + # Инициируем генерацию заголовка + 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" + yield f"data: {json.dumps({'type': 'done', 'node_id': assistant_node_id})}\n\n" except Exception as e: print(f"Ошибка при стриминге ответа: {e}") - error_data = {'type': 'error', 'error': str(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') @@ -172,31 +211,41 @@ def chat_stream(): @api.route('/api/chat/regenerate', methods=['POST']) def regenerate_message(): - """API endpoint для регенерации сообщения.""" + """API endpoint для регенерации сообщения. + + При ошибке LLM-узла, переиспользует его. В противном случае, создает новый узел. + """ data = request.get_json() graph_id = data.get("graph_id") - node_id = data.get("node_id") - '''system_prompt = data.get("system_prompt") - model = data.get("model")''' + node_id = data.get("node_id") # Это ID узла, который нужно регенерировать (llm-ответ) + model = data.get("model") if not graph_id or not node_id: return jsonify({"error": "Не указаны graph_id или node_id."}, 400) try: - # Определяем родительский узел - parent_node_id = graph_history_manager.get_parent_node_id( - graph_id, node_id) + # Определяем родительский узел для 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) + return jsonify({"error": "Не удалось найти родительский узел для регенерации."}, 404) - # Создаём новый узел-заглушку для регенерации - new_assistant_node_id = graph_history_manager.create_assistant_placeholder_node( - graph_id, parent_node_id) + # НОВОЕ: Проверяем тип узла + 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({ - "new_node_id": new_assistant_node_id, - "parent_node_id": parent_node_id + "existing_assistant_node_id": assistant_node_id_to_use, # Возвращаем ID узла, который будет регенерирован/использован + "user_node_id": parent_node_id, # Это родительский узел (пользовательский) + "graph_id": graph_id, + "model": model # Передаем выбранную модель }) except Exception as e: print(f"Ошибка при регенерации сообщения: {e}") diff --git a/app/graph_history_manager.py b/app/graph_history_manager.py index 983d33b..7873aca 100644 --- a/app/graph_history_manager.py +++ b/app/graph_history_manager.py @@ -500,6 +500,48 @@ class GraphHistoryManager: print(f"Ошибка при создании узла-заглушки: {e}") raise + + def mark_node_as_placeholder_and_clear_content(self, graph_id: str, node_id: str, parent_node_id: str): + """ + Очищает содержимое существующего узла LLM, помечает его как заглушку, + и обновляет current_node_id графа. + @param graph_id: ID графа. + @param node_id: ID узла, который нужно очистить и пометить как заглушку. + @param parent_node_id: ID родительского узла (пользовательского), к которому будет прикреплен этот LLM-узел. + """ + with self._get_connection() as conn: + cursor = conn.cursor() + try: + # 1. Получаем текущие данные узла + cursor.execute("SELECT node_data FROM graph_nodes_data WHERE graph_id = ? AND node_id = ?", + (graph_id, node_id)) + result = cursor.fetchone() + + if result: + node_data = json.loads(result[0]) + # Очищаем контент и сбрасываем состояние + node_data["message"]["content"] = "" + node_data["label"] = "Получение ответа..." + node_data["is_placeholder"] = True + # Сбрасываем флаги генерации заголовка для перегенерации + node_data["title"] = None + node_data["title_generated"] = False + + # Обновляем узел + cursor.execute("UPDATE graph_nodes_data SET node_type = ?, node_data = ?, title = NULL, title_generated = FALSE WHERE graph_id = ? AND node_id = ?", + ("llm", json.dumps(node_data), graph_id, node_id)) + + # 2. Обновляем current_node_id графа на этот узел + self._update_graph_current_node_id(conn, graph_id, node_id) + + conn.commit() + else: + raise ValueError(f"Узел с ID {node_id} не найден в графе {graph_id} для обновления.") + except Exception as e: + conn.rollback() + print(f"Ошибка при очистке и маркировке узла {node_id} как заглушки: {e}") + raise + def update_assistant_node_content(self, graph_id: str, node_id: str, content: str): """ Обновляет содержимое узла ассистента после завершения стриминга. @@ -560,3 +602,17 @@ class GraphHistoryManager: (graph_id, node_id)) result = cursor.fetchone() return result[0] if result else None + + def get_node_type(self, graph_id: str, node_id: str) -> Optional[str]: + """ + Возвращает тип узла по его ID. + @param graph_id: ID графа. + @param node_id: ID узла, тип которого нужно получить. + @returns: Тип узла (например, 'user', 'llm', 'error') или None, если узел не найден. + """ + with self._get_connection() as conn: + cursor = conn.cursor() + cursor.execute("SELECT node_type FROM graph_nodes_data WHERE graph_id = ? AND node_id = ?", + (graph_id, node_id)) + result = cursor.fetchone() + return result[0] if result else None \ No newline at end of file diff --git a/app/llm_client.py b/app/llm_client.py index 79a25da..b44e867 100644 --- a/app/llm_client.py +++ b/app/llm_client.py @@ -84,6 +84,16 @@ MODELS: Dict[str, Dict[str, Any]] = { "stream": True, "capabilities": ["vision"], }, + "mistral-small-latest-error": { + "name": "mistral-small-latest", + "provider": "mistralai", + "model_name": "mistral-small-latest-error", # Добавлено имя модели для LangChain + "apiBase": + "https://render-service-gsu7.onrender.com/m", + "apiKey": "Q0m29fvxBY0Cfdj4sjHaKqccy1NjonLW", + "stream": True, + "capabilities": ["vision"], + }, } diff --git a/app/nodes.py b/app/nodes.py index 3dcdd78..22cd8dc 100644 --- a/app/nodes.py +++ b/app/nodes.py @@ -143,7 +143,7 @@ def execute_command_node(state: AgentState) -> AgentState: DEFAULT_LLM_NAME = "gemini-2.0-flash" #"gemini-2.5-flash" main_llm = get_llm(DEFAULT_LLM_NAME) -DEFAULT_SUMMARIZATION_LLM_NAME = "gemini-2.0-flash-r" # "mistral-small-latest" #"gemini-2.0-flash-r" +DEFAULT_SUMMARIZATION_LLM_NAME = "mistral-small-latest" # "mistral-small-latest" #"gemini-2.0-flash-r" def call_llm_node(state: AgentState) -> AgentState: diff --git a/app/title_generator.py b/app/title_generator.py index bf45f72..899e2f0 100644 --- a/app/title_generator.py +++ b/app/title_generator.py @@ -25,7 +25,7 @@ class TitleGenerator: self.thread: Optional[threading.Thread] = None self.socketio = None - + # Системные промпты для генерации заголовков self.graph_title_prompt = """Создай краткий заголовок (максимум 80 символов) для диалога на основе первого сообщения пользователя. Заголовок должен отражать основную тему или вопрос. Заголовок должен быть простым текстом, без какого-либо форматирования или использования @@ -199,11 +199,11 @@ class TitleGenerator: # ----------------------------------------------- WebSocket Notification --------------------------------------------------------------- # Отправляем событие через WebSocket if self.socketio: - self.socketio.emit('graph_title_updated', { + self._emit_with_retry('graph_title_updated', { 'graph_id': graph_id, 'title': title - }) - + }, max_retries=5) + except Exception as e: print(f"Ошибка генерации заголовка графа {graph_id}: {e}") @@ -248,17 +248,17 @@ class TitleGenerator: # Сохраняем заголовок self.history_manager.update_node_title(graph_id, node_id, title) - print(f"Сгенерирован заголовок узла {node_id}: {title}") + print(f"✅ Сгенерирован заголовок узла {node_id}: `{title}` из текста `{content[:40]}...`") # ----------------------------------------------- WebSocket Notification --------------------------------------------------------------- # Отправляем событие через WebSocket if self.socketio: - self.socketio.emit('node_title_updated', { + self._emit_with_retry('node_title_updated', { 'graph_id': graph_id, 'node_id': node_id, 'title': title - }) - + }, max_retries=5) + except Exception as e: print(f"Ошибка генерации заголовка узла {node_id}: {e}") @@ -289,4 +289,52 @@ class TitleGenerator: cursor.execute( "INSERT INTO title_generation_queue (item_type, graph_id, priority) VALUES (?, ?, ?)", ("graph", graph_id, priority)) - conn.commit() \ No newline at end of file + conn.commit() + + def _emit_with_retry(self, event_name: str, data: dict, max_retries: int = 5): + """ + Отправляет WebSocket событие с механизмом повторных попыток и подтверждением. + @param event_name: Название события. + @param data: Данные для отправки. + @param max_retries: Максимальное количество попыток отправки. + """ + if not self.socketio: + return + + retry_count = 0 + ack_received = threading.Event() + + def ack_callback(response): + """Колбэк, вызываемый при получении подтверждения от клиента.""" + if response and response.get('status') == 'ok': + ack_received.set() + print(f"✅ Получено подтверждение для {event_name}: node_id={data.get('node_id')}/graph_id={data.get('graph_id')}") + else: + print(f"⚠️ Получен некорректный ответ для {event_name}: {response}") + + while retry_count < max_retries and not ack_received.is_set(): + retry_count += 1 + + try: + print(f"📤 Попытка {retry_count}/{max_retries} отправки {event_name} для node_id={data.get('node_id')}/graph_id={data.get('graph_id')}") + + # Отправляем событие с callback для подтверждения + self.socketio.emit(event_name, data, callback=ack_callback) + + # Ждем подтверждения до 2 секунд + if ack_received.wait(timeout=2.0): + return # Успешно получено подтверждение + + print(f"⏱️ Таймаут ожидания подтверждения для {event_name} (попытка {retry_count})") + + # Небольшая пауза перед повторной попыткой + if retry_count < max_retries: + time.sleep(0.5) + + except Exception as e: + print(f"❌ Ошибка при отправке {event_name} (попытка {retry_count}): {e}") + if retry_count < max_retries: + time.sleep(0.5) + + if not ack_received.is_set(): + print(f"❌ Не удалось доставить {event_name} после {max_retries} попыток: {data}") \ No newline at end of file