"""
Асинхронный сервис для генерации заголовков графов и узлов.
Работает в фоновом режиме, обрабатывая очередь элементов без заголовков.
"""
import sqlite3
import threading
import time
from typing import Any, Dict, Optional
from llm_client import get_llm
from langchain_core.messages import SystemMessage, HumanMessage
import os
DEFAULT_SUMMARIZATION_LLM_NAME = "gemini-2.0-flash-lite" # "mistral-small-latest" # "mistral-small-latest" #"gemini-2.0-flash-r"
DEFAULT_VOICE_LLM_NAME = "gemini-2.0-flash-lite" # "gemini-3.0-flash-openrouter"
MAX_TITLE_GENERATION_CONTENT_LENGTH = 5000 # Максимальное количество символов для генерации заголовков
from app.voice_service import VOICE_COMMANDS_RESPONSE_TO_STORE
class TitleGenerator:
"""Сервис для асинхронной генерации заголовков."""
def __init__(self, history_manager, db_path="graph_history.db"):
self.db_path = db_path
self._create_table()
self.history_manager = history_manager
self.llm = get_llm(DEFAULT_SUMMARIZATION_LLM_NAME)
self.voice_llm = get_llm(DEFAULT_VOICE_LLM_NAME)
self.running = False
self.thread: Optional[threading.Thread] = None
self.socketio = None
# Системные промпты для генерации заголовков
self.graph_title_prompt = """Создай краткий заголовок (максимум 80 символов) для диалога на основе первого сообщения пользователя.
Заголовок должен отражать основную тему или вопрос. Заголовок должен быть простым текстом, без какого-либо форматирования или использования
специальных символов разметки (например, #, *, _, `). Отвечай только заголовком, без дополнительных объяснений. Ты должен сделать саммери, а не ответить на вопросы, если они есть в сообщении."""
self.node_title_prompt = """Создай краткий заголовок (максимум 40 символов) для этого сообщения/действия.
Заголовок должен кратко описывать суть сообщения или действия. Заголовок должен быть простым текстом, без какого-либо форматирования или использования
специальных символов разметки (например, #, *, _, `). Отвечай только заголовком, без дополнительных объяснений. Ты должен сделать саммери, а не ответить на вопросы, если они есть в сообщении."""
def set_socketio(self, socketio):
"""
Устанавливает объект SocketIO для отправки уведомлений.
@param socketio: Объект SocketIO.
"""
self.socketio = socketio
def _get_connection(self):
"""
Получает соединение с базой данных.
Устанавливаем таймаут для ожидания блокировки базы данных при конкурентной записи.
"""
return sqlite3.connect(self.db_path,
timeout=5.0) # Увеличил таймаут до 5 секунд
def _create_table(self):
"""Создает таблицы для хранения графов, если они не существуют."""
with self._get_connection() as conn:
cursor = conn.cursor()
# Очередь для обработки заголовков
cursor.execute("""
CREATE TABLE IF NOT EXISTS title_generation_queue (
id INTEGER PRIMARY KEY AUTOINCREMENT,
item_type TEXT NOT NULL, -- 'graph' or 'node'
graph_id TEXT NOT NULL,
node_id TEXT DEFAULT NULL, -- только для узлов
priority INTEGER DEFAULT 0, -- для приоритета обработки
created_at DATETIME DEFAULT CURRENT_TIMESTAMP
)
""")
conn.commit()
def start(self):
"""Запускает фоновый процесс генерации заголовков."""
if self.running:
return
self.running = True
# Заполняем очередь при старте
self.populate_initial_title_queue()
# Запускаем фоновый поток
self.thread = threading.Thread(target=self._process_queue, daemon=True)
self.thread.start()
print("Сервис генерации заголовков запущен")
def populate_initial_title_queue(self):
"""Заполняет очередь элементами без сгенерированных заголовков при старте сервера."""
with self._get_connection() as conn:
cursor = conn.cursor()
# Очищаем существующую очередь
cursor.execute("DELETE FROM title_generation_queue")
# Добавляем графы без заголовков
cursor.execute(
"SELECT id FROM graphs WHERE title_generated = FALSE OR title_generated IS NULL"
)
for (graph_id, ) in cursor.fetchall():
cursor.execute(
"INSERT INTO title_generation_queue (item_type, graph_id, priority) VALUES (?, ?, ?)",
("graph", graph_id, 1))
# Добавляем узлы без заголовков
cursor.execute(
"SELECT graph_id, node_id FROM graph_nodes_data WHERE title_generated = FALSE OR title_generated IS NULL"
)
for graph_id, node_id in cursor.fetchall():
cursor.execute(
"INSERT INTO title_generation_queue (item_type, graph_id, node_id, priority) VALUES (?, ?, ?, ?)",
("node", graph_id, node_id, 1))
conn.commit()
def stop(self):
"""Останавливает фоновый процесс."""
self.running = False
if self.thread:
self.thread.join()
print("Сервис генерации заголовков остановлен")
def _process_queue(self):
"""Основной цикл обработки очереди заголовков."""
while self.running:
try:
item = self.get_next_from_title_queue()
if item:
if item["item_type"] == "graph":
self._generate_graph_title(item["graph_id"])
elif item["item_type"] == "node":
self._generate_node_title(item["graph_id"], item["node_id"])
# ОБРАБОТКА ГОЛОСА
elif item["item_type"] == "voice_command":
self._process_voice_command(item["node_id"], item["graph_id"])
elif item["item_type"] == "voice_regular":
self._process_voice_regular(item["node_id"], item["graph_id"])
else:
# Если очередь пуста, ждем немного
time.sleep(5)
except Exception as e:
print(f"Ошибка при обработке очереди заголовков: {e}")
time.sleep(10)
def get_next_from_title_queue(self) -> Optional[Dict[str, Any]]:
"""Получает следующий элемент из очереди для обработки."""
with self._get_connection() as conn:
cursor = conn.cursor()
cursor.execute("""
SELECT id, item_type, graph_id, node_id
FROM title_generation_queue
ORDER BY priority DESC, id ASC
LIMIT 1
""")
result = cursor.fetchone()
if result:
queue_id, item_type, graph_id, node_id = result
# Удаляем из очереди
cursor.execute(
"DELETE FROM title_generation_queue WHERE id = ?",
(queue_id, ))
conn.commit()
return {
"item_type": item_type,
"graph_id": graph_id,
"node_id": node_id
}
return None
def _generate_graph_title(self, graph_id: str):
"""Генерирует заголовок для графа."""
try:
graph_data = self.history_manager.get_graph(graph_id)
if not graph_data:
return
messages = graph_data.get("messages", [])
if not messages:
return
# Находим первое пользовательское сообщение
first_user_message = None
for msg in messages:
if msg.get("role") == "user":
# Ограничиваем сообщение первыми ... символами для генерации заголовка
first_user_message = msg.get("content", "")[:MAX_TITLE_GENERATION_CONTENT_LENGTH]
break
if not first_user_message:
return
# Генерируем заголовок
llm_messages = [
SystemMessage(content=self.graph_title_prompt),
HumanMessage(content=first_user_message)
]
title = self.llm.invoke(llm_messages)
title = title.strip()[:80] # Ограничиваем длину
# Сохраняем заголовок
self.history_manager.update_graph_title(graph_id, title)
print(f"Сгенерирован заголовок графа {graph_id}: {title}")
# ----------------------------------------------- WebSocket Notification ---------------------------------------------------------------
# Отправляем событие через WebSocket
if self.socketio:
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}")
def _generate_node_title(self, graph_id: str, node_id: str):
"""Генерирует заголовок для узла."""
try:
graph_data = self.history_manager.get_graph(graph_id)
if not graph_data:
return
# Находим узел по ID
target_node = None
for node in graph_data.get("graph_nodes", []):
if node["id"] == node_id:
target_node = node
break
if not target_node:
return
# Получаем содержимое узла
node_data = target_node.get("data", {})
message = node_data.get("message", {})
# Ограничиваем сообщение первыми ... символами для генерации заголовка
content = message.get("content", "")[:MAX_TITLE_GENERATION_CONTENT_LENGTH]
if not content:
# Если нет content, используем label или type
content = node_data.get("label", target_node.get("type", ""))
if not content:
return
# Генерируем заголовок
llm_messages = [
SystemMessage(content=self.node_title_prompt),
HumanMessage(content=content)
]
title = self.llm.invoke(llm_messages)
title = title.strip()[:40] # Ограничиваем длину
# Сохраняем заголовок
self.history_manager.update_node_title(graph_id, node_id, title)
print(f"✅ Сгенерирован заголовок узла {node_id}: `{title}` из текста `{content[:40]}...`")
# ----------------------------------------------- WebSocket Notification ---------------------------------------------------------------
# Отправляем событие через WebSocket
if self.socketio:
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}")
def add_node_to_queue(self, cursor, graph_id, node_id):
# Добавляем в очередь заголовков с высоким приоритетом
cursor.execute(
"INSERT INTO title_generation_queue (item_type, graph_id, node_id, priority) VALUES (?, ?, ?, ?)",
("node", graph_id, node_id, 10))
def add_graph_to_queue(self, cursor, graph_id):
cursor.execute(
"INSERT INTO title_generation_queue (item_type, graph_id, priority) VALUES (?, ?, ?)",
("graph", graph_id, 10))
def add_node_to_queue_direct(self, graph_id: str, node_id: str, priority: int = 10):
"""Добавляет узел в очередь генерации заголовков напрямую (с собственным соединением)."""
with self._get_connection() as conn:
cursor = conn.cursor()
cursor.execute(
"INSERT INTO title_generation_queue (item_type, graph_id, node_id, priority) VALUES (?, ?, ?, ?)",
("node", graph_id, node_id, priority))
conn.commit()
def add_graph_to_queue_direct(self, graph_id: str, priority: int = 10):
"""Добавляет граф в очередь генерации заголовков напрямую (с собственным соединением)."""
with self._get_connection() as conn:
cursor = conn.cursor()
cursor.execute(
"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}")
# VOICE
def add_voice_task_to_queue(self, command_text, prompt, priority=20):
"""Добавляет задачу на обработку голосовой команды в очередь."""
with self._get_connection() as conn:
cursor = conn.cursor()
# Используем node_id для хранения текста команды,
# а item_type 'voice_command' для идентификации
cursor.execute(
"INSERT INTO title_generation_queue (item_type, graph_id, node_id, priority) VALUES (?, ?, ?, ?)",
("voice_command", prompt, command_text, priority))
conn.commit()
# 2. Добавьте метод постановки регулярной задачи
def add_voice_regular_task(self, transcription_text, prompt, priority=15):
with self._get_connection() as conn:
cursor = conn.cursor()
cursor.execute(
"INSERT INTO title_generation_queue (item_type, graph_id, node_id, priority) VALUES (?, ?, ?, ?)",
("voice_regular", prompt, transcription_text, priority))
conn.commit()
# 3. Добавьте метод обработки регулярной задачи
def _process_voice_regular(self, text, prompt):
try:
print(f"🎤 Анализ транскрибации (Prompt 1)...")
llm_messages = [
SystemMessage(content=prompt),
HumanMessage(content=f"Последний фрагмент диалога для анализа: {text}")
]
response = self.voice_llm.invoke(llm_messages)
import app.api as api
api.voice_inst.responses["regular"] = response # Записываем в Response 1
print(f"✅🎤 Регулярный анализ голосового ввода завершен.")
except Exception as e:
print(f"❌🎤 Ошибка регулярного анализа: {e}")
# 4. Поправьте _process_voice_command (запись команды в Response 2)
def _process_voice_command(self, text, prompt):
try:
print(f"🎤 Обработка голосовой команды (Prompt 2)...")
llm_messages = [
SystemMessage(content=prompt),
HumanMessage(content=f"Выполни команду из транскрибации: {text}")
]
response = self.voice_llm.invoke(llm_messages)
import app.api as api
# Форматируем для вкладки Response 2
formatted_res = f"Команда: {text}
Ответ: {response}"
api.voice_inst.responses["commands"].insert(0, formatted_res)
# Ограничиваем историю команд, чтобы не раздувать память
if len(api.voice_inst.responses["commands"]) > VOICE_COMMANDS_RESPONSE_TO_STORE:
api.voice_inst.responses["commands"].pop(0)
print(f"✅🎤 Голосовая команда обработана.")
except Exception as e:
print(f"❌🎤 Ошибка обработки команды: {e}")