More correct working with error nodes + more correct regeneration

This commit is contained in:
dimitrievgs 2025-10-12 21:40:50 +03:00
parent 21c9e70de6
commit 196201c681
6 changed files with 199 additions and 36 deletions

2
.gitignore vendored
View File

@ -6,4 +6,4 @@ node_modules/
__pycache__
*.db
*.db-history
*.db-journal

View File

@ -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}")

View File

@ -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

View File

@ -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"],
},
}

View File

@ -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:

View File

@ -199,10 +199,10 @@ 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,16 +248,16 @@ 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}")
@ -290,3 +290,51 @@ class TitleGenerator:
"INSERT INTO title_generation_queue (item_type, graph_id, priority) VALUES (?, ?, ?)",
("graph", graph_id, priority))
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}")