diff --git a/AsyncTask.md b/AsyncTask.md new file mode 100644 index 0000000..9ee94a9 --- /dev/null +++ b/AsyncTask.md @@ -0,0 +1,196 @@ +Отличная идея! Гарантировать уникальность `id` узлов и ребер является критически важным для целостности графа. SQLite с `PRIMARY KEY` и `INSERT OR IGNORE` прекрасно с этим справляется, автоматически пропуская попытки вставить уже существующие записи. Если нам не нужна строгая последовательность `id`, а только их уникальность, то это работает очень хорошо. + +Что касается `_get_max_graph_id`, ты прав. Если мы создаём `id` как `graph_1`, `graph_2` и так далее, то `_get_max_graph_id` работает как способ получения **следующего свободного числового суффикса** для нового графа. При этом совершенно неважно, в каком порядке графы были *добавлены* или *изменены* в БД. Важно, что каждый новый граф получит уникальный ID. + +### Детальное ТЗ для изменений + +**Цель:** Обеспечить целостность и конкурентную работу с графами в многопоточной среде Flask, используя SQLite. Предотвратить состояние гонки, гарантировать уникальность ID узлов/ребер и атомарность операций. + +**Основные изменения:** + +1. **Изменение модели хранения графа в БД:** + * Вместо хранения `graph_nodes` и `graph_edges` как JSON-строк в одной строке таблицы `graphs`, будет создано 3 таблицы: + * `graphs`: Основная таблица для информации о графе (`id`, `current_node_id`). + * `graph_nodes_data`: Для хранения каждого узла графа как отдельной записи (`graph_id`, `node_id`, `node_type`, `node_data_json`). + * `graph_edges_data`: Для хранения каждого ребра графа как отдельной записи (`graph_id`, `edge_id`, `source_node_id`, `target_node_id`). + * Это позволит инкрементально добавлять узлы и ребра без перезаписи всего графа целиком. + +2. **Обработка конкуренции:** + * **Удаление `threading.Lock()`:** Мы не будем использовать мьютекс Python для блокировки всего `GraphHistoryManager`. Вместо этого мы будем полагаться на встроенные механизмы конкурентной обработки SQLite и атомарности транзакций. + * **Транзакции:** Все операции по модификации графа (добавление узлов, ребер, обновление `current_node_id`) будут обернуты в одну SQLite транзакцию. Это гарантирует, что либо все изменения будут применены, либо ни одно. + * **`INSERT OR IGNORE`:** В операциях добавления узлов и ребер будет использоваться `INSERT OR IGNORE`. Это важно: если два конкурирующих запроса попытаются добавить один и тот же новый узел/ребро (с одинаковым `id`), первый запрос успешно его добавит, а второй будет проигнорирован без ошибки. Это предотвращает "Invalid ID" ошибки при инкрементальном добавлении. + * **Timeout для подключения SQLite:** Увеличение таймаута при подключении к SQLite (`timeout=5.0`) даст другим потокам больше времени на освобождение блокировки файлов БД при пиковых нагрузках, избегая `sqlite3.OperationalError: database is locked`. + +3. **Генерация ID:** + * `_get_max_graph_id` будет использоваться только для генерации нового уникального `graph_id`, если он не предоставлен вызывающим кодом (т.е. создание нового графа). + * В `nodes.py` генерация `node_id` будет опираться на `UUID`, чтобы гарантировать уникальность без необходимости глобального отслеживания. Это критично, поскольку разные потоки могут одновременно генерировать узлы для разных веток. + +4. **`run_agent` логика:** + * `run_agent` будет загружать полный граф (все узлы и ребра) из `GraphHistoryManager` в начале. + * После выполнения LangGraph, `run_agent` будет сравнивать `final_state.graph_nodes` и `final_state.graph_edges` с исходными загруженными данными, чтобы определить *только новые* узлы и ребра, которые нужно сохранить. + * Затем `graph_history_manager.save_graph_changes` будет вызван с этим списком новых узлов/ребер и новым `current_node_id`. + +**Файлы для изменения:** + +1. `app/graph_history_manager.py` +2. `app/workflows.py` +3. `app/nodes.py` + +--- + +Хорошо, полностью согласен с переходом на UUID для генерации ID графов, узлов и ребер. Это упрощает логику, исключает необходимость в `_get_max_graph_id` и гарантирует глобальную уникальность ID без дополнительных сложностей с нумерацией. + +### Детальное ТЗ для изменений (Обновленное) + +**Цель:** Обеспечить целостность и конкурентную работу с графами в многопоточной среде Flask, используя SQLite. Предотвратить состояние гонки, гарантировать уникальность ID узлов/ребер через UUID и атомарность операций. + +**Основные изменения:** + +1. **Генерация ID с использованием UUID:** + * Все ID для графов (`graph_id`), узлов (`node_id`) и ребер (`edge_id`) будут генерироваться с использованием `uuid.uuid4().hex`. + * Это исключает необходимость в методе `_get_max_graph_id`. + +2. **Изменение модели хранения графа в БД:** (Остается как в предыдущем ТЗ) + * 3 таблицы: `graphs`, `graph_nodes_data`, `graph_edges_data`. + * Инкрементальное добавление узлов и ребер без перезаписи всего графа. + +3. **Обработка конкуренции:** (Остается как в предыдущем ТЗ) + * Удаление `threading.Lock()` из `GraphHistoryManager`. + * Все операции по модификации графа в `save_graph_changes` будут обернуты в одну SQLite транзакцию. + * Использование `INSERT OR IGNORE` для узлов и ребер. + * Timeout при подключении к SQLite (`timeout=5.0`). + +4. **`run_agent` логика:** (Остается как в предыдущем ТЗ) + * `run_agent` будет загружать полный граф (все узлы и ребра) из `GraphHistoryManager` в начале. + * После выполнения LangGraph, `run_agent` будет сравнивать `final_state.graph_nodes` и `final_state.graph_edges` с исходными загруженными данными, чтобы определить *только новые* узлы и ребра, которые нужно сохранить. + * Затем `graph_history_manager.save_graph_changes` будет вызван с этим списком новых узлов/ребер и новым `current_node_id`. + +**Файлы для изменения:** + +1. `app/graph_history_manager.py` +2. `app/workflows.py` +3. `app/nodes.py` + +--- + +### ТЗ для `app/graph_history_manager.py` + +**Цель:** Адаптировать менеджер истории графов для работы с UUID и обеспечить надежное инкрементальное хранение и конкурентный доступ. + +**Изменения:** + +1. **Импорты:** + * Удалить `threading`. + +2. **Удаление `_get_max_graph_id`:** + * Полностью удалить этот метод, так как ID будут генерироваться с помощью UUID. + +3. **Метод `save_graph_changes`:** + * **Генерация `graph_id`**: Если `existing_graph_id` равен `None`, новый `graph_id` генерируется внутри метода с использованием `uuid.uuid4().hex`. + * Пример: `graph_id = uuid.uuid4().hex` + * Логика сохранения остальных изменений остаётся прежней (использование `INSERT OR IGNORE`, транзакции). + * Важно: В `INSERT OR IGNORE INTO graphs (id) VALUES (?)` теперь просто используем сгенерированный или переданный `graph_id`. + +4. **Методы `_add_graph_edge` и `_add_graph_node`:** + * Эти методы должны использовать `INSERT OR IGNORE` для добавления узлов и ребер. Это гарантирует, что если узел или ребро с таким `(graph_id, node_id)` или `(graph_id, edge_id)` уже существует из-за конкурентной операции, попытка вставки будет проигнорирована без ошибки. + * Предполагается, что `node["id"]` и `edge["id"]` уже должны быть UUID, сгенерированными на этапе создания узла/ребра. + * **Внимание для `_add_graph_edge`**: SQL-запрос `FOREIGN KEY` требует, чтобы узлы `source_node_id` и `target_node_id` уже существовали. Если для `INSERT OR IGNORE` возникнет нарушение `FOREIGN KEY`, он все равно выдаст ошибку `IntegrityError`. Нужно убедиться, что узлы всегда добавляются *до* ребер, которые на них ссылаются, что обеспечивается логикой в `workflows.py`. Если мы будем использовать `executescript` как в предыдущем примере, то нужно убедиться, что синтаксис верен для `INSERT OR IGNORE`. Простой `cursor.execute` с `INSERT OR IGNORE` также должен работать, если узлы гарантированно существуют. + + Пример `_add_graph_edge` с `cursor.execute`: + ```python + def _add_graph_edge(self, conn, graph_id: str, edge: Dict[str, Any]): + cursor = conn.cursor() + cursor.execute( + "INSERT OR IGNORE INTO graph_edges_data (graph_id, edge_id, source_node_id, target_node_id) VALUES (?, ?, ?, ?)", + (graph_id, edge["id"], edge["source"], edge["target"]) + ) + ``` + Этот вариант предпочтительнее `executescript`, так как он более читаем и безопасен от SQL-инъекций. + +--- + +### ТЗ для `app/workflows.py` + +**Цель:** Адаптировать логику запуска агента для работы с UUID и инкрементального сохранения изменений в графе. + +**Изменения:** + +1. **Импорты:** + * Добавить `import uuid`. + * Удалить `GraphUpdateConflictError` из импортов, так как этот exception больше не используется прямой логикой `run_agent`. + +2. **Метод `run_agent`:** + * **Загрузка графа:** `loaded_graph_data` будет содержать `graph_nodes` и `graph_edges`. + * **Отслеживание исходного состояния:** После загрузки графа, перед вызовом `app.invoke`, необходимо сохранить текущие `graph_nodes` и `graph_edges` из `initial_state` (которые были загружены из БД). Это нужно, чтобы потом определить, какие узлы/ребра являются *новыми*. + * Пример: + ```python + original_nodes_ids = {node['id'] for node in initial_state.graph_nodes} + original_edges_ids = {edge['id'] for edge in initial_state.graph_edges} + ``` + * **Генерация `graph_id` при создании нового графа:** Если `existing_graph_id` равен `None`, то `graph_id` должен быть сгенерирован здесь, до `app.invoke`, и передан в `initial_state`. + * Пример: + ```python + if existing_graph_id is None: + new_graph_id = uuid.uuid4().hex + # Также может быть полезно добавить его в initial_state, если AgentState это поддерживает + # initial_state.graph_id = new_graph_id + ``` + (Или просто передавать `None` в `save_graph_changes` и пусть генерация там происходит). + * **Определение новых узлов и ребер для сохранения:** После получения `final_state` из `app.invoke`: + * Проитерировать `final_state.graph_nodes` и собрать те, `id` которых нет в `original_nodes_ids`. Это будут `newly_added_nodes`. + * Аналогично для `final_state.graph_edges` и `original_edges_ids`, чтобы получить `newly_added_edges`. + * **Вызов `graph_history_manager.save_graph_changes`:** + * Вызвать `save_graph_changes` с `existing_graph_id` (или сгенерированным `new_graph_id` если он новый), `newly_added_nodes`, `newly_added_edges` и `final_state.current_node_id`. + * Пример: + ```python + final_graph_id = existing_graph_id or new_graph_id # если new_graph_id был сгенерирован ранее + graph_id_after_save = graph_history_manager.save_graph_changes( + final_graph_id, newly_added_nodes, newly_added_edges, final_state.current_node_id + ) + ``` + * **Обновление ID в `response_data`**: Убедиться, что `response_data["graph_id"]` содержит актуальный `graph_id_after_save`. + +--- + +### ТЗ для `app/nodes.py` + +**Цель:** Модифицировать генерацию ID узлов и ребер на UUID. + +**Изменения:** + +1. **Импорты:** + * Добавить `import uuid`. + +2. **Генерация `node_id`:** + * Во всех функциях узлов (`parse_command_node`, `call_llm_node`, `generate_images_node`, `analyze_image_node`, `get_meet_subtitles_node`, `get_teams_subtitles_node`, `summarize_history_node`, `handle_error_node`, `help_node`) заменить строки типа `f"llm_{len(state.graph_nodes) + 1}"` на `uuid.uuid4().hex`. + * Это обеспечит уникальность `node_id` каждого нового узла. + +3. **Генерация `edge_id`:** + * При создании ребер (например, `f"e{state.parent_node_id}-{user_node_id}"`) также заменять на `uuid.uuid4().hex`. + * Пример: + ```python + edge_id = uuid.uuid4().hex + state.graph_edges.append({ + "id": edge_id, + "source": state.parent_node_id, + "target": user_node_id + }) + ``` + +--- + +**По поводу фронтенда:** + +**Да, фронтенд, скорее всего, придется менять.** + +Причина в том, что теперь бэкенд ожидает, что вы можете отправить `parent_node_id` в запросах к `/api/chat`. Если ваш фронтенд не был разработан с учетом "продолжения" чата с определенного узла (т.е. просто отправлял `message` и `graph_id`), то теперь ему нужно будет: + +1. **Отслеживать `current_node_id`:** После каждого ответа от бэкенда (`response.graph_visualization_data.current_node_id` или `response.messages[-1].node_id`), фронтенду нужно будет сохранять этот `node_id`. +2. **Отправлять `parent_node_id`:** Когда пользователь хочет продолжить чат с какого-либо узла (например, нажав на него или просто отправляя новый запрос после получения ответа), фронтенд должен будет включать этот `current_node_id` (или `node_id` того узла, с которого он хочет продолжить) в качестве `parent_node_id` в Payload POST-запроса к `/api/chat`. + +Если фронтенд отправляет только `message` и `graph_id`, и вы хотите простой линейный диалог, где новый ответ всегда продолжается с последнего сгенерированного узла, то `parent_node_id` можно не отправлять, и бэкенд будет продолжать с `current_node_id`, который он хранит для `graph_id`. Однако, если вы хотите реализовать возможность "ответить на конкретное сообщение/узел", то `parent_node_id` необходим. + +**Резюме:** + +* **`api.py`:** Нужны изменения, как показано выше. +* **Фронтенд:** Потенциально требует изменений для использования новой функциональности `parent_node_id` и более гибкого управления историей чата/графом. \ No newline at end of file diff --git a/app/graph_history_manager.py b/app/graph_history_manager.py index 6ec81c0..416d0ba 100644 --- a/app/graph_history_manager.py +++ b/app/graph_history_manager.py @@ -1,135 +1,183 @@ -""" -Этот модуль управляет сохранением и извлечением истории графов -запросов, используя SQLite базу данных для хранения данных. -""" - +# отключаем 40-ка строчное ограничение для этого файла +# ----------------------------------------------- Imports ------------------------------------------------------------ from typing import Dict, Any, List, Optional - import json - import sqlite3 +import uuid # Для генерации UUID +# ----------------------------------------------- Exceptions ------------------------------------------------------------ +class GraphUpdateConflictError(Exception): + """Исключение возникает при конфликте обновления графа.""" + pass +# ----------------------------------------------- GraphHistoryManager ------------------------------------------------------------ class GraphHistoryManager: """ Менеджер для хранения истории графов запросов в SQLite. + Использует отдельные таблицы для узлов и ребер для инкрементального обновления. + Обеспечивает конкурентный доступ к разным графам, полагаясь на транзакционность SQLite. """ def __init__(self, db_path="graph_history.db"): self.db_path = db_path - self._create_table(self._get_connection()) # Создаем таблицу при инициализации + self._create_tables() + # ----------------------------------------------- Internal Methods ------------------------------------------------------------ def _get_connection(self): - """Получает соединение с базой данных.""" - return sqlite3.connect(self.db_path) + """ + Получает соединение с базой данных. + Устанавливаем таймаут для ожидания блокировки базы данных при конкурентной записи. + """ + return sqlite3.connect(self.db_path, timeout=5.0) # Увеличил таймаут до 5 секунд - def _create_table(self, conn): - """Создает таблицу для хранения графов, если она не существует.""" + def _create_tables(self): + """Создает таблицы для хранения графов, если они не существуют.""" + with self._get_connection() as conn: + cursor = conn.cursor() + cursor.execute(""" + CREATE TABLE IF NOT EXISTS graphs ( + id TEXT PRIMARY KEY, + current_node_id TEXT DEFAULT NULL, + timestamp DATETIME DEFAULT CURRENT_TIMESTAMP -- Добавляем метку времени для сортировки + ) + """) + cursor.execute(""" + CREATE TABLE IF NOT EXISTS graph_nodes_data ( + graph_id TEXT NOT NULL, + node_id TEXT NOT NULL, + node_type TEXT NOT NULL, + node_data TEXT NOT NULL, -- JSON-строка для 'data' из узла + PRIMARY KEY (graph_id, node_id), + FOREIGN KEY (graph_id) REFERENCES graphs (id) ON DELETE CASCADE + ) + """) + cursor.execute(""" + CREATE TABLE IF NOT EXISTS graph_edges_data ( + graph_id TEXT NOT NULL, + edge_id TEXT NOT NULL, + source_node_id TEXT NOT NULL, + target_node_id TEXT NOT NULL, + PRIMARY KEY (graph_id, edge_id), + FOREIGN KEY (graph_id) REFERENCES graphs (id) ON DELETE CASCADE, + FOREIGN KEY (source_node_id) REFERENCES graph_nodes_data (node_id) ON DELETE CASCADE, + FOREIGN KEY (target_node_id) REFERENCES graph_nodes_data (node_id) ON DELETE CASCADE + ) + """) + conn.commit() + + # _get_max_graph_id - удален, так как используем UUID для graph_id + + # ----------------------------------------------- Graph Node/Edge Access Helpers ------------------------------------------------------------ + def _get_graph_nodes(self, conn, graph_id: str) -> List[Dict[str, Any]]: + """Извлекает все узлы для заданного графа.""" cursor = conn.cursor() - cursor.execute(""" - CREATE TABLE IF NOT EXISTS graphs ( - id TEXT PRIMARY KEY, - graph_nodes TEXT, - graph_edges TEXT, - current_node_id TEXT - ) - """) - conn.commit() + cursor.execute("SELECT node_id, node_type, node_data FROM graph_nodes_data WHERE graph_id = ?", (graph_id,)) + nodes = [] + for node_id, node_type, node_data_json in cursor.fetchall(): + nodes.append({ + "id": node_id, + "type": node_type, + "data": json.loads(node_data_json) + }) + return nodes - def _get_max_graph_id(self, conn): - """Получает максимальный ID графа из базы данных.""" + def _get_graph_edges(self, conn, graph_id: str) -> List[Dict[str, Any]]: + """Извлекает все ребра для заданного графа.""" + cursor = conn.cursor() + cursor.execute("SELECT edge_id, source_node_id, target_node_id FROM graph_edges_data WHERE graph_id = ?", (graph_id,)) + edges = [] + for edge_id, source, target in cursor.fetchall(): + edges.append({ + "id": edge_id, + "source": source, + "target": target + }) + return edges + + def _update_graph_current_node_id(self, conn, graph_id: str, new_current_node_id: str): + """Обновляет current_node_id для графа.""" + cursor = conn.cursor() + cursor.execute("UPDATE graphs SET current_node_id = ? WHERE id = ?", (new_current_node_id, graph_id)) + + def _add_graph_node(self, conn, graph_id: str, node: Dict[str, Any]): + """Добавляет новый узел в граф. Использует INSERT OR IGNORE для избежания конфликтов.""" cursor = conn.cursor() cursor.execute( - "SELECT MAX(CAST(SUBSTR(id, 7) AS INTEGER)) FROM graphs WHERE id LIKE 'graph_%'" + "INSERT OR IGNORE INTO graph_nodes_data (graph_id, node_id, node_type, node_data) VALUES (?, ?, ?, ?)", + (graph_id, node["id"], node["type"], json.dumps(node["data"])) ) - result = cursor.fetchone()[0] - return result if result is not None else 0 - # ---------------------------------------------------- public methods ---------------------------------------------------------------- - def save_graph(self, graph_data: Any) -> str: - """ - Сохраняет или обновляет данные графа в базе данных и возвращает ID. - Не сохраняет полную историю сообщений, только структуру графа. - """ - conn = self._get_connection() + def _add_graph_edge(self, conn, graph_id: str, edge: Dict[str, Any]): + """Добавляет новое ребро в граф. Использует INSERT OR IGNORE для избежания конфликтов.""" cursor = conn.cursor() - try: - graph_id = graph_data.get("id") - if not graph_id: - max_graph_id = self._get_max_graph_id(conn) - next_graph_id = max_graph_id + 1 - graph_id = f"graph_{next_graph_id}" - graph_data["id"] = graph_id + cursor.execute( + "INSERT OR IGNORE INTO graph_edges_data (graph_id, edge_id, source_node_id, target_node_id) VALUES (?, ?, ?, ?)", + (graph_id, edge["id"], edge["source"], edge["target"]) + ) - # Преобразуем структуры данных в JSON-строки для хранения в SQLite - graph_nodes_json = json.dumps(graph_data.get("graph_nodes", [])) - graph_edges_json = json.dumps(graph_data.get("graph_edges", [])) - current_node_id = graph_data.get("current_node_id", "") + # ----------------------------------------------- Public Methods ------------------------------------------------------------ + def save_graph_changes(self, graph_id: str, + new_nodes: List[Dict[str, Any]], + new_edges: List[Dict[str, Any]], + current_node_id: str) -> str: + """ + Сохраняет *изменения* в граф в базе данных. + Добавляет новые узлы и ребра и обновляет current_node_id. + Операции проводятся в рамках одной транзакции. + Конфликты при добавлении существующих узлов/ребер игнорируются. + Если graph_id пуст, генерируется новый UUID для графа. + """ + if not graph_id: + graph_id = str(uuid.uuid4()) # Генерируем новый UUID для ID графа + + with self._get_connection() as conn: # Одна транзакция для всех изменений + try: + # Убедимся, что запись о графе существует в главной таблице `graphs` + # INSERT OR IGNORE создаст новую запись, если её нет. + # Если граф уже существует, это ничего не изменит. + conn.execute("INSERT OR IGNORE INTO graphs (id) VALUES (?)", (graph_id,)) + + # Добавляем новые узлы + for node in new_nodes: + self._add_graph_node(conn, graph_id, node) - # Проверяем, существует ли уже запись с таким ID - cursor.execute("SELECT id FROM graphs WHERE id = ?", (graph_id, )) - existing_record = cursor.fetchone() + # Добавляем новые ребра + for edge in new_edges: + self._add_graph_edge(conn, graph_id, edge) - if existing_record: - # Обновляем существующую запись - cursor.execute( - """ - UPDATE graphs SET - graph_nodes = ?, - graph_edges = ?, - current_node_id = ? - WHERE id = ? - """, (graph_nodes_json, graph_edges_json, - current_node_id, graph_id)) - else: - # Вставляем новую запись - cursor.execute( - """ - INSERT INTO graphs (id, graph_nodes, graph_edges, current_node_id) - VALUES (?, ?, ?, ?) - """, (graph_id, graph_nodes_json, - graph_edges_json, current_node_id)) + # Обновляем current_node_id - это атомарное изменение в рамках транзакции + self._update_graph_current_node_id(conn, graph_id, current_node_id) - conn.commit() - print(f"Граф сохранен/обновлен: {graph_id}") - return graph_id - finally: - conn.close() + conn.commit() # Фиксируем все изменения + print(f"Изменения графа {graph_id} сохранены.") + return graph_id + except sqlite3.OperationalError as e: + conn.rollback() # Откатываем транзакцию при ошибке (например, DB Locked) + print(f"Ошибка блокировки SQLite при сохранении изменений графа {graph_id}: {e}") + raise + except Exception as e: + conn.rollback() # Откатываем транзакцию при других ошибках + print(f"Ошибка при сохранении изменений графа {graph_id}: {e}") + raise def get_graph(self, graph_id: str, target_node_id: Optional[str] = None) -> Optional[Dict[str, Any]]: """ Получает данные графа по ID из базы данных. - Динамически вычисляет 'messages' до 'target_node_id'. - Если target_node_id не указан, используется current_node_id из БД. """ - conn = self._get_connection() - cursor = conn.cursor() - try: - cursor.execute( - "SELECT graph_nodes, graph_edges, current_node_id FROM graphs WHERE id = ?", - (graph_id, )) + with self._get_connection() as conn: + cursor = conn.cursor() + cursor.execute("SELECT current_node_id FROM graphs WHERE id = ?", (graph_id,)) result = cursor.fetchone() if result: - graph_nodes_json, graph_edges_json, db_current_node_id = result + db_current_node_id = result[0] - # Преобразуем JSON-строки обратно в структуры данных Python - graph_nodes = json.loads(graph_nodes_json) - graph_edges = json.loads(graph_edges_json) + # Получаем все узлы и ребра + graph_nodes = self._get_graph_nodes(conn, graph_id) + graph_edges = self._get_graph_edges(conn, graph_id) resolved_current_node_id = target_node_id if target_node_id else db_current_node_id - # Если нет узлов, значит граф пуст или некорректен - if not graph_nodes: - return { - "id": graph_id, - "messages": [], - "graph_nodes": [], - "graph_edges": [], - "current_node_id": resolved_current_node_id - } - - # Динамически собираем сообщения messages = self.get_messages_from_root_to_node( {"graph_nodes": graph_nodes, "graph_edges": graph_edges}, resolved_current_node_id @@ -137,31 +185,27 @@ class GraphHistoryManager: return { "id": graph_id, - "messages": messages, # Здесь будут вычисленные сообщения - "graph_nodes": graph_nodes, - "graph_edges": graph_edges, + "messages": messages, + "graph_nodes": graph_nodes, + "graph_edges": graph_edges, "current_node_id": resolved_current_node_id } else: return None - finally: - conn.close() def get_all_graphs_summary(self) -> List[Dict[str, str]]: """ - Возвращает краткий список всех сохраненных графов. - Первое сообщение извлекается путем построения пути к current_node_id - и взятия первого сообщения. + Возвращает краткий список всех сохраненных графов, отсортированных по дате добавления. """ - conn = self._get_connection() - cursor = conn.cursor() - try: - cursor.execute("SELECT id, graph_nodes, graph_edges, current_node_id FROM graphs") + with self._get_connection() as conn: + cursor = conn.cursor() + # Сортируем по timestamp по убыванию, чтобы самые новые были сверху + cursor.execute("SELECT id, current_node_id FROM graphs ORDER BY timestamp DESC") graphs_data = cursor.fetchall() summaries = [] - for gid, nodes_json, edges_json, current_node_id in graphs_data: - graph_nodes = json.loads(nodes_json) - graph_edges = json.loads(edges_json) + for gid, current_node_id in graphs_data: + graph_nodes = self._get_graph_nodes(conn, gid) + graph_edges = self._get_graph_edges(conn, gid) # Получаем все сообщения до current_node_id full_messages = self.get_messages_from_root_to_node( @@ -171,21 +215,18 @@ class GraphHistoryManager: first_message_content = "No content" if full_messages: - # Ищем первое сообщение, которое является пользовательским for msg in full_messages: if msg.get("role") == "user": first_message_content = msg.get("content", "No content") break - if first_message_content == "No content": # Если не нашли пользовательского, берем первое ассистента - first_message_content = full_messages[0].get("content", "No content") + if first_message_content == "No content" and full_messages: + first_message_content = full_messages[0].get("content", "No content") summaries.append({ "id": gid, "first_message": first_message_content }) return summaries - finally: - conn.close() def get_messages_from_root_to_node( self, graph_data: Dict[str, Any], @@ -202,19 +243,15 @@ class GraphHistoryManager: # Создаем словарь для быстрого поиска родительских узлов parent_map = {edge['target']: edge['source'] for edge in edges} - # Функция для рекурсивного подъема по дереву до корневого узла def get_path_to_root(node_id: str) -> List[str]: path = [node_id] current = node_id while current in parent_map: current = parent_map[current] path.append(current) - return path[::-1] # Инвертируем, чтобы получить путь от корня + return path[::-1] - # Получаем путь от корня до целевого узла path_to_target = get_path_to_root(target_node_id) - - # Собираем сообщения, соответствующие узлам в пути. # Сообщения теперь хранятся в 'data' каждого узла при создании. messages_for_path = [] @@ -222,7 +259,6 @@ class GraphHistoryManager: node = node_map.get(node_id) if node and 'message' in node.get('data', {}): # Проверяем наличие поля 'message' original_message = node['data']['message'] - # Извлекаем 'role', 'content' и добавляем 'type' из node['type'] messages_for_path.append({ "role": original_message.get("role", "unknown"), "type": node.get("type", "unknown"), # Используем 'type' узла @@ -232,16 +268,19 @@ class GraphHistoryManager: return messages_for_path + def delete_graph(self, graph_id: str) -> bool: - """Удаляет граф из базы данных.""" - conn = self._get_connection() - cursor = conn.cursor() - try: - cursor.execute("DELETE FROM graphs WHERE id = ?", (graph_id, )) - conn.commit() - return cursor.rowcount > 0 # Возвращает True, если была удалена хотя бы одна строка - except Exception as e: - print(f"Ошибка при удалении графа: {e}") - return False - finally: - conn.close() + """ + Удаляет граф из базы данных. + Благодаря CASCADE, удалит также все узлы и ребра, связанные с этим графом. + """ + with self._get_connection() as conn: + cursor = conn.cursor() + try: + cursor.execute("DELETE FROM graphs WHERE id = ?", (graph_id,)) + conn.commit() + return cursor.rowcount > 0 + except Exception as e: + conn.rollback() + print(f"Ошибка при удалении графа: {e}") + return False \ No newline at end of file diff --git a/app/nodes.py b/app/nodes.py index 3f694e4..3d6c988 100644 --- a/app/nodes.py +++ b/app/nodes.py @@ -11,6 +11,7 @@ from services import ImageGenerationService, SubtitleService, ImageAnalysisServi from models import AgentState from typing import Dict, Any, Optional +import uuid # Добавлено для генерации UUID class CommandManager: @@ -74,8 +75,8 @@ def parse_command_node(state: AgentState) -> AgentState: print("Выполняется узел: parse_command_node") user_input = state.input.strip() - # Создаем ID для узла пользователя - user_node_id = f"user_{len(state.graph_nodes) + 1}" + # Создаем ID для узла пользователя с помощью UUID + user_node_id = str(uuid.uuid4()) # Добавляем узел пользователя в граф для фронтенда с полным текстом сообщения state.graph_nodes.append({ @@ -89,8 +90,9 @@ def parse_command_node(state: AgentState) -> AgentState: # Добавляем ребро к новому узлу, если есть родительский узел if state.parent_node_id: + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{user_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": user_node_id }) @@ -160,7 +162,7 @@ def call_llm_node(state: AgentState) -> AgentState: response = main_llm.invoke(messages_for_llm) - llm_node_id = f"llm_{len(state.graph_nodes) + 1}" + llm_node_id = str(uuid.uuid4()) # Добавляем узел LLM в граф для фронтенда с полным текстом ответа state.graph_nodes.append({ @@ -171,8 +173,9 @@ def call_llm_node(state: AgentState) -> AgentState: "message": {"role": "assistant", "content": response, "node_id": llm_node_id} # Сохраняем сообщение здесь } }) + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{llm_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": llm_node_id }) @@ -224,7 +227,7 @@ def generate_images_node(state: AgentState) -> AgentState: response_text = f"Сгенерировано {len(urls)} изображений:\n{image_markdown}" state.llm_response = response_text - img_gen_node_id = f"img_gen_{len(state.graph_nodes) + 1}" + img_gen_node_id = str(uuid.uuid4()) # Добавляем узел генерации изображений в граф state.graph_nodes.append({ "id": img_gen_node_id, @@ -235,8 +238,9 @@ def generate_images_node(state: AgentState) -> AgentState: "message": {"role": "assistant", "content": response_text, "node_id": img_gen_node_id} } }) + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{img_gen_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": img_gen_node_id }) @@ -285,7 +289,7 @@ def analyze_image_node(state: AgentState) -> AgentState: response_text = f"Результат анализа изображения '{image_url}': {analysis_result}" state.llm_response = response_text - img_analysis_node_id = f"img_analysis_{len(state.graph_nodes) + 1}" + img_analysis_node_id = str(uuid.uuid4()) # Добавляем узел анализа изображения в граф state.graph_nodes.append({ "id": img_analysis_node_id, @@ -296,8 +300,9 @@ def analyze_image_node(state: AgentState) -> AgentState: "message": {"role": "assistant", "content": response_text, "node_id": img_analysis_node_id} } }) + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{img_analysis_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": img_analysis_node_id }) @@ -334,7 +339,7 @@ def get_meet_subtitles_node(state: AgentState) -> AgentState: response_text = f"Текущие субтитры из Google Meet: {subtitles}" state.llm_response = response_text - meet_subtitles_node_id = f"meet_subtitles_{len(state.graph_nodes) + 1}" + meet_subtitles_node_id = str(uuid.uuid4()) # Добавляем узел субтитров Meet в граф state.graph_nodes.append({ "id": meet_subtitles_node_id, @@ -345,8 +350,9 @@ def get_meet_subtitles_node(state: AgentState) -> AgentState: "message": {"role": "assistant", "content": response_text, "node_id": meet_subtitles_node_id} } }) + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{meet_subtitles_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": meet_subtitles_node_id }) @@ -380,7 +386,7 @@ def get_teams_subtitles_node(state: AgentState) -> AgentState: response_text = f"Текущие субтитры из MS Teams: {subtitles}" state.llm_response = response_text - teams_subtitles_node_id = f"teams_subtitles_{len(state.graph_nodes) + 1}" + teams_subtitles_node_id = str(uuid.uuid4()) # Добавляем узел субтитров Teams в граф state.graph_nodes.append({ "id": teams_subtitles_node_id, @@ -391,8 +397,9 @@ def get_teams_subtitles_node(state: AgentState) -> AgentState: "message": {"role": "assistant", "content": response_text, "node_id": teams_subtitles_node_id} } }) + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{teams_subtitles_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": teams_subtitles_node_id }) @@ -432,7 +439,7 @@ def summarize_history_node(state: AgentState) -> AgentState: response_text = f"История диалога суммирована:\n{summarized_text}" state.llm_response = response_text - summary_node_id = f"summary_{len(state.graph_nodes) + 1}" + summary_node_id = str(uuid.uuid4()) # Добавляем узел суммаризации в граф state.graph_nodes.append({ "id": summary_node_id, @@ -443,8 +450,9 @@ def summarize_history_node(state: AgentState) -> AgentState: "message": {"role": "assistant", "content": response_text, "node_id": summary_node_id} } }) + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{summary_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": summary_node_id }) @@ -474,7 +482,7 @@ def handle_error_node(state: AgentState) -> AgentState: print(f"Выполняется узел: handle_error_node. Ошибка: {state.error}") state.llm_response = state.error - error_node_id = f"error_{len(state.graph_nodes) + 1}" + error_node_id = str(uuid.uuid4()) # Добавляем узел ошибки в граф state.graph_nodes.append({ "id": error_node_id, @@ -485,8 +493,9 @@ def handle_error_node(state: AgentState) -> AgentState: "message": {"role": "assistant", "content": state.llm_response, "node_id": error_node_id} } }) + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{error_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": error_node_id }) @@ -514,7 +523,7 @@ def help_node(state: AgentState) -> AgentState: help_text += f"/{cmd}: {info['description']}\n" state.llm_response = help_text - help_node_id = f"help_{len(state.graph_nodes) + 1}" + help_node_id = str(uuid.uuid4()) # Добавляем узел справки в граф state.graph_nodes.append({ "id": help_node_id, @@ -525,8 +534,9 @@ def help_node(state: AgentState) -> AgentState: "message": {"role": "assistant", "content": state.llm_response, "node_id": help_node_id} } }) + edge_id = str(uuid.uuid4()) state.graph_edges.append({ - "id": f"e{state.parent_node_id}-{help_node_id}", + "id": edge_id, "source": state.parent_node_id, "target": help_node_id }) @@ -540,4 +550,4 @@ def help_node(state: AgentState) -> AgentState: "node_id": help_node_id }) - return state + return state \ No newline at end of file diff --git a/app/workflows.py b/app/workflows.py index 46e5339..4a8745f 100644 --- a/app/workflows.py +++ b/app/workflows.py @@ -4,11 +4,12 @@ и переходами между узлами обработки. """ -from typing import Dict, Any, Optional +from typing import Dict, Any, Optional, Set from langgraph.graph import StateGraph, END from graph_history_manager import GraphHistoryManager from models import AgentState from nodes import parse_command_node, execute_command_node, call_llm_node, generate_images_node, analyze_image_node, get_meet_subtitles_node, get_teams_subtitles_node, summarize_history_node, handle_error_node, help_node +import uuid # Добавлено для генерации UUID # --- LangGraph: Построение графа --- graph_history_manager = GraphHistoryManager() @@ -100,6 +101,10 @@ def run_agent(user_input: str, # ---------------------------------------------------- Initial State Setup ---------------------------------------------------------------- initial_state = AgentState(input=user_input, parent_node_id=parent_node_id) + + # Сохраняем исходные ID узлов и ребер для определения новых после выполнения графа + original_node_ids: Set[str] = set() + original_edge_ids: Set[str] = set() # Если есть существующий ID графа, загружаем его историю и структуру loaded_graph_data = None @@ -115,6 +120,10 @@ def run_agent(user_input: str, initial_state.graph_nodes = loaded_graph_data.get("graph_nodes", []) initial_state.graph_edges = loaded_graph_data.get("graph_edges", []) + # Сохраняем ID загруженных узлов и ребер + original_node_ids = {node['id'] for node in initial_state.graph_nodes} + original_edge_ids = {edge['id'] for edge in initial_state.graph_edges} + # Устанавливаем parent_node_id для нового входного узла # Если parent_node_id был передан в запросе, используем его, # иначе берем current_node_id из загруженного графа. @@ -123,6 +132,10 @@ def run_agent(user_input: str, # История чата для текущего раунда работы LLM формируется из загруженных сообщений. initial_state.temporary_chat_history = loaded_graph_data.get("messages", []) + # Определяем graph_id, который будет использоваться для сохранения. + # Если существующий ID не передан, GraphHistoryManager сгенерирует новый. + graph_id_to_save = existing_graph_id + # ---------------------------------------------------- run_agent ---------------------------------------------------------------- # Запускаем граф @@ -130,23 +143,25 @@ def run_agent(user_input: str, final_state = AgentState(**result) # ---------------------------------------------------- Save Graph ---------------------------------------------------------------- + # Определяем новые узлы и ребра для сохранения + newly_added_nodes = [node for node in final_state.graph_nodes if node['id'] not in original_node_ids] + newly_added_edges = [edge for edge in final_state.graph_edges if edge['id'] not in original_edge_ids] + # Сохраняем обновленное состояние графа. # Обратите внимание: chat_history больше не сохраняется напрямую, # она является частью данных узлов в graph_nodes. - graph_data_to_save = { - "id": existing_graph_id, # Используем существующий ID или он будет сгенерирован - # менеджером истории при первом сохранении - "graph_nodes": final_state.graph_nodes, - "graph_edges": final_state.graph_edges, - "current_node_id": final_state.current_node_id - } - new_graph_id = graph_history_manager.save_graph(graph_data_to_save) + final_graph_id = graph_history_manager.save_graph_changes( + graph_id_to_save, + newly_added_nodes, + newly_added_edges, + final_state.current_node_id + ) # ---------------------------------------------------- prepare response ---------------------------------------------------------------- # Для `messages` в `response_data` используем `final_state.temporary_chat_history`, # так как она отражает только сообщения текущей ветки, добавленные в ходе этого выполнения. response_data = { - "graph_id": new_graph_id, + "graph_id": final_graph_id, "messages": final_state.temporary_chat_history, "llm_response": final_state.llm_response, "image_urls": final_state.image_urls, @@ -160,4 +175,4 @@ def run_agent(user_input: str, "current_node_id": final_state.current_node_id } } - return response_data + return response_data \ No newline at end of file diff --git a/graph_history.db b/graph_history.db index 45037c0..e820d3e 100644 Binary files a/graph_history.db and b/graph_history.db differ