changing graph state via atomic changes instead of changing whole graph, ids via uuid

This commit is contained in:
dimitrievgs 2025-09-16 00:14:14 +03:00
parent 51affba1a5
commit bb4ba2644c
5 changed files with 423 additions and 163 deletions

196
AsyncTask.md Normal file
View File

@ -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` и более гибкого управления историей чата/графом.

View File

@ -1,135 +1,183 @@
""" # отключаем 40-ка строчное ограничение для этого файла
Этот модуль управляет сохранением и извлечением истории графов # ----------------------------------------------- Imports ------------------------------------------------------------
запросов, используя SQLite базу данных для хранения данных.
"""
from typing import Dict, Any, List, Optional from typing import Dict, Any, List, Optional
import json import json
import sqlite3 import sqlite3
import uuid # Для генерации UUID
# ----------------------------------------------- Exceptions ------------------------------------------------------------
class GraphUpdateConflictError(Exception):
"""Исключение возникает при конфликте обновления графа."""
pass
# ----------------------------------------------- GraphHistoryManager ------------------------------------------------------------
class GraphHistoryManager: class GraphHistoryManager:
""" """
Менеджер для хранения истории графов запросов в SQLite. Менеджер для хранения истории графов запросов в SQLite.
Использует отдельные таблицы для узлов и ребер для инкрементального обновления.
Обеспечивает конкурентный доступ к разным графам, полагаясь на транзакционность SQLite.
""" """
def __init__(self, db_path="graph_history.db"): def __init__(self, db_path="graph_history.db"):
self.db_path = db_path self.db_path = db_path
self._create_table(self._get_connection()) # Создаем таблицу при инициализации self._create_tables()
# ----------------------------------------------- Internal Methods ------------------------------------------------------------
def _get_connection(self): 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 = conn.cursor()
cursor.execute(""" cursor.execute("""
CREATE TABLE IF NOT EXISTS graphs ( CREATE TABLE IF NOT EXISTS graphs (
id TEXT PRIMARY KEY, id TEXT PRIMARY KEY,
graph_nodes TEXT, current_node_id TEXT DEFAULT NULL,
graph_edges TEXT, timestamp DATETIME DEFAULT CURRENT_TIMESTAMP -- Добавляем метку времени для сортировки
current_node_id TEXT )
""")
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() conn.commit()
def _get_max_graph_id(self, conn): # _get_max_graph_id - удален, так как используем UUID для graph_id
"""Получает максимальный ID графа из базы данных."""
# ----------------------------------------------- Graph Node/Edge Access Helpers ------------------------------------------------------------
def _get_graph_nodes(self, conn, graph_id: str) -> List[Dict[str, Any]]:
"""Извлекает все узлы для заданного графа."""
cursor = conn.cursor()
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_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 = conn.cursor()
cursor.execute( 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 _add_graph_edge(self, conn, graph_id: str, edge: Dict[str, Any]):
def save_graph(self, graph_data: Any) -> str: """Добавляет новое ребро в граф. Использует INSERT OR IGNORE для избежания конфликтов."""
"""
Сохраняет или обновляет данные графа в базе данных и возвращает ID.
Не сохраняет полную историю сообщений, только структуру графа.
"""
conn = self._get_connection()
cursor = conn.cursor() cursor = conn.cursor()
try: cursor.execute(
graph_id = graph_data.get("id") "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"])
)
# ----------------------------------------------- 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: if not graph_id:
max_graph_id = self._get_max_graph_id(conn) graph_id = str(uuid.uuid4()) # Генерируем новый UUID для ID графа
next_graph_id = max_graph_id + 1
graph_id = f"graph_{next_graph_id}"
graph_data["id"] = graph_id
# Преобразуем структуры данных в JSON-строки для хранения в SQLite with self._get_connection() as conn: # Одна транзакция для всех изменений
graph_nodes_json = json.dumps(graph_data.get("graph_nodes", [])) try:
graph_edges_json = json.dumps(graph_data.get("graph_edges", [])) # Убедимся, что запись о графе существует в главной таблице `graphs`
current_node_id = graph_data.get("current_node_id", "") # INSERT OR IGNORE создаст новую запись, если её нет.
# Если граф уже существует, это ничего не изменит.
conn.execute("INSERT OR IGNORE INTO graphs (id) VALUES (?)", (graph_id,))
# Проверяем, существует ли уже запись с таким ID # Добавляем новые узлы
cursor.execute("SELECT id FROM graphs WHERE id = ?", (graph_id, )) for node in new_nodes:
existing_record = cursor.fetchone() self._add_graph_node(conn, graph_id, node)
if existing_record: # Добавляем новые ребра
# Обновляем существующую запись for edge in new_edges:
cursor.execute( self._add_graph_edge(conn, graph_id, edge)
"""
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))
conn.commit() # Обновляем current_node_id - это атомарное изменение в рамках транзакции
print(f"Граф сохранен/обновлен: {graph_id}") self._update_graph_current_node_id(conn, graph_id, current_node_id)
conn.commit() # Фиксируем все изменения
print(f"Изменения графа {graph_id} сохранены.")
return graph_id return graph_id
finally: except sqlite3.OperationalError as e:
conn.close() 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]]: def get_graph(self, graph_id: str, target_node_id: Optional[str] = None) -> Optional[Dict[str, Any]]:
""" """
Получает данные графа по ID из базы данных. Получает данные графа по ID из базы данных.
Динамически вычисляет 'messages' до 'target_node_id'.
Если target_node_id не указан, используется current_node_id из БД.
""" """
conn = self._get_connection() with self._get_connection() as conn:
cursor = conn.cursor() cursor = conn.cursor()
try: cursor.execute("SELECT current_node_id FROM graphs WHERE id = ?", (graph_id,))
cursor.execute(
"SELECT graph_nodes, graph_edges, current_node_id FROM graphs WHERE id = ?",
(graph_id, ))
result = cursor.fetchone() result = cursor.fetchone()
if result: 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_nodes = self._get_graph_nodes(conn, graph_id)
graph_edges = json.loads(graph_edges_json) 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 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( messages = self.get_messages_from_root_to_node(
{"graph_nodes": graph_nodes, "graph_edges": graph_edges}, {"graph_nodes": graph_nodes, "graph_edges": graph_edges},
resolved_current_node_id resolved_current_node_id
@ -137,31 +185,27 @@ class GraphHistoryManager:
return { return {
"id": graph_id, "id": graph_id,
"messages": messages, # Здесь будут вычисленные сообщения "messages": messages,
"graph_nodes": graph_nodes, "graph_nodes": graph_nodes,
"graph_edges": graph_edges, "graph_edges": graph_edges,
"current_node_id": resolved_current_node_id "current_node_id": resolved_current_node_id
} }
else: else:
return None return None
finally:
conn.close()
def get_all_graphs_summary(self) -> List[Dict[str, str]]: def get_all_graphs_summary(self) -> List[Dict[str, str]]:
""" """
Возвращает краткий список всех сохраненных графов. Возвращает краткий список всех сохраненных графов, отсортированных по дате добавления.
Первое сообщение извлекается путем построения пути к current_node_id
и взятия первого сообщения.
""" """
conn = self._get_connection() with self._get_connection() as conn:
cursor = conn.cursor() cursor = conn.cursor()
try: # Сортируем по timestamp по убыванию, чтобы самые новые были сверху
cursor.execute("SELECT id, graph_nodes, graph_edges, current_node_id FROM graphs") cursor.execute("SELECT id, current_node_id FROM graphs ORDER BY timestamp DESC")
graphs_data = cursor.fetchall() graphs_data = cursor.fetchall()
summaries = [] summaries = []
for gid, nodes_json, edges_json, current_node_id in graphs_data: for gid, current_node_id in graphs_data:
graph_nodes = json.loads(nodes_json) graph_nodes = self._get_graph_nodes(conn, gid)
graph_edges = json.loads(edges_json) graph_edges = self._get_graph_edges(conn, gid)
# Получаем все сообщения до current_node_id # Получаем все сообщения до current_node_id
full_messages = self.get_messages_from_root_to_node( full_messages = self.get_messages_from_root_to_node(
@ -171,12 +215,11 @@ class GraphHistoryManager:
first_message_content = "No content" first_message_content = "No content"
if full_messages: if full_messages:
# Ищем первое сообщение, которое является пользовательским
for msg in full_messages: for msg in full_messages:
if msg.get("role") == "user": if msg.get("role") == "user":
first_message_content = msg.get("content", "No content") first_message_content = msg.get("content", "No content")
break break
if first_message_content == "No content": # Если не нашли пользовательского, берем первое ассистента if first_message_content == "No content" and full_messages:
first_message_content = full_messages[0].get("content", "No content") first_message_content = full_messages[0].get("content", "No content")
summaries.append({ summaries.append({
@ -184,8 +227,6 @@ class GraphHistoryManager:
"first_message": first_message_content "first_message": first_message_content
}) })
return summaries return summaries
finally:
conn.close()
def get_messages_from_root_to_node( def get_messages_from_root_to_node(
self, graph_data: Dict[str, Any], self, graph_data: Dict[str, Any],
@ -202,19 +243,15 @@ class GraphHistoryManager:
# Создаем словарь для быстрого поиска родительских узлов # Создаем словарь для быстрого поиска родительских узлов
parent_map = {edge['target']: edge['source'] for edge in edges} parent_map = {edge['target']: edge['source'] for edge in edges}
# Функция для рекурсивного подъема по дереву до корневого узла
def get_path_to_root(node_id: str) -> List[str]: def get_path_to_root(node_id: str) -> List[str]:
path = [node_id] path = [node_id]
current = node_id current = node_id
while current in parent_map: while current in parent_map:
current = parent_map[current] current = parent_map[current]
path.append(current) path.append(current)
return path[::-1] # Инвертируем, чтобы получить путь от корня return path[::-1]
# Получаем путь от корня до целевого узла
path_to_target = get_path_to_root(target_node_id) path_to_target = get_path_to_root(target_node_id)
# Собираем сообщения, соответствующие узлам в пути. # Собираем сообщения, соответствующие узлам в пути.
# Сообщения теперь хранятся в 'data' каждого узла при создании. # Сообщения теперь хранятся в 'data' каждого узла при создании.
messages_for_path = [] messages_for_path = []
@ -222,7 +259,6 @@ class GraphHistoryManager:
node = node_map.get(node_id) node = node_map.get(node_id)
if node and 'message' in node.get('data', {}): # Проверяем наличие поля 'message' if node and 'message' in node.get('data', {}): # Проверяем наличие поля 'message'
original_message = node['data']['message'] original_message = node['data']['message']
# Извлекаем 'role', 'content' и добавляем 'type' из node['type']
messages_for_path.append({ messages_for_path.append({
"role": original_message.get("role", "unknown"), "role": original_message.get("role", "unknown"),
"type": node.get("type", "unknown"), # Используем 'type' узла "type": node.get("type", "unknown"), # Используем 'type' узла
@ -232,16 +268,19 @@ class GraphHistoryManager:
return messages_for_path return messages_for_path
def delete_graph(self, graph_id: str) -> bool: def delete_graph(self, graph_id: str) -> bool:
"""Удаляет граф из базы данных.""" """
conn = self._get_connection() Удаляет граф из базы данных.
Благодаря CASCADE, удалит также все узлы и ребра, связанные с этим графом.
"""
with self._get_connection() as conn:
cursor = conn.cursor() cursor = conn.cursor()
try: try:
cursor.execute("DELETE FROM graphs WHERE id = ?", (graph_id, )) cursor.execute("DELETE FROM graphs WHERE id = ?", (graph_id,))
conn.commit() conn.commit()
return cursor.rowcount > 0 # Возвращает True, если была удалена хотя бы одна строка return cursor.rowcount > 0
except Exception as e: except Exception as e:
conn.rollback()
print(f"Ошибка при удалении графа: {e}") print(f"Ошибка при удалении графа: {e}")
return False return False
finally:
conn.close()

View File

@ -11,6 +11,7 @@ from services import ImageGenerationService, SubtitleService, ImageAnalysisServi
from models import AgentState from models import AgentState
from typing import Dict, Any, Optional from typing import Dict, Any, Optional
import uuid # Добавлено для генерации UUID
class CommandManager: class CommandManager:
@ -74,8 +75,8 @@ def parse_command_node(state: AgentState) -> AgentState:
print("Выполняется узел: parse_command_node") print("Выполняется узел: parse_command_node")
user_input = state.input.strip() user_input = state.input.strip()
# Создаем ID для узла пользователя # Создаем ID для узла пользователя с помощью UUID
user_node_id = f"user_{len(state.graph_nodes) + 1}" user_node_id = str(uuid.uuid4())
# Добавляем узел пользователя в граф для фронтенда с полным текстом сообщения # Добавляем узел пользователя в граф для фронтенда с полным текстом сообщения
state.graph_nodes.append({ state.graph_nodes.append({
@ -89,8 +90,9 @@ def parse_command_node(state: AgentState) -> AgentState:
# Добавляем ребро к новому узлу, если есть родительский узел # Добавляем ребро к новому узлу, если есть родительский узел
if state.parent_node_id: if state.parent_node_id:
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{user_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": user_node_id "target": user_node_id
}) })
@ -160,7 +162,7 @@ def call_llm_node(state: AgentState) -> AgentState:
response = main_llm.invoke(messages_for_llm) response = main_llm.invoke(messages_for_llm)
llm_node_id = f"llm_{len(state.graph_nodes) + 1}" llm_node_id = str(uuid.uuid4())
# Добавляем узел LLM в граф для фронтенда с полным текстом ответа # Добавляем узел LLM в граф для фронтенда с полным текстом ответа
state.graph_nodes.append({ 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} # Сохраняем сообщение здесь "message": {"role": "assistant", "content": response, "node_id": llm_node_id} # Сохраняем сообщение здесь
} }
}) })
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{llm_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": llm_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}" response_text = f"Сгенерировано {len(urls)} изображений:\n{image_markdown}"
state.llm_response = response_text 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({ state.graph_nodes.append({
"id": img_gen_node_id, "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} "message": {"role": "assistant", "content": response_text, "node_id": img_gen_node_id}
} }
}) })
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{img_gen_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": img_gen_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}" response_text = f"Результат анализа изображения '{image_url}': {analysis_result}"
state.llm_response = response_text 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({ state.graph_nodes.append({
"id": img_analysis_node_id, "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} "message": {"role": "assistant", "content": response_text, "node_id": img_analysis_node_id}
} }
}) })
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{img_analysis_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": img_analysis_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}" response_text = f"Текущие субтитры из Google Meet: {subtitles}"
state.llm_response = response_text 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 в граф # Добавляем узел субтитров Meet в граф
state.graph_nodes.append({ state.graph_nodes.append({
"id": meet_subtitles_node_id, "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} "message": {"role": "assistant", "content": response_text, "node_id": meet_subtitles_node_id}
} }
}) })
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{meet_subtitles_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": meet_subtitles_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}" response_text = f"Текущие субтитры из MS Teams: {subtitles}"
state.llm_response = response_text 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 в граф # Добавляем узел субтитров Teams в граф
state.graph_nodes.append({ state.graph_nodes.append({
"id": teams_subtitles_node_id, "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} "message": {"role": "assistant", "content": response_text, "node_id": teams_subtitles_node_id}
} }
}) })
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{teams_subtitles_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": teams_subtitles_node_id "target": teams_subtitles_node_id
}) })
@ -432,7 +439,7 @@ def summarize_history_node(state: AgentState) -> AgentState:
response_text = f"История диалога суммирована:\n{summarized_text}" response_text = f"История диалога суммирована:\n{summarized_text}"
state.llm_response = response_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({ state.graph_nodes.append({
"id": summary_node_id, "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} "message": {"role": "assistant", "content": response_text, "node_id": summary_node_id}
} }
}) })
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{summary_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": summary_node_id "target": summary_node_id
}) })
@ -474,7 +482,7 @@ def handle_error_node(state: AgentState) -> AgentState:
print(f"Выполняется узел: handle_error_node. Ошибка: {state.error}") print(f"Выполняется узел: handle_error_node. Ошибка: {state.error}")
state.llm_response = 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({ state.graph_nodes.append({
"id": error_node_id, "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} "message": {"role": "assistant", "content": state.llm_response, "node_id": error_node_id}
} }
}) })
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{error_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": error_node_id "target": error_node_id
}) })
@ -514,7 +523,7 @@ def help_node(state: AgentState) -> AgentState:
help_text += f"/{cmd}: {info['description']}\n" help_text += f"/{cmd}: {info['description']}\n"
state.llm_response = help_text 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({ state.graph_nodes.append({
"id": help_node_id, "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} "message": {"role": "assistant", "content": state.llm_response, "node_id": help_node_id}
} }
}) })
edge_id = str(uuid.uuid4())
state.graph_edges.append({ state.graph_edges.append({
"id": f"e{state.parent_node_id}-{help_node_id}", "id": edge_id,
"source": state.parent_node_id, "source": state.parent_node_id,
"target": help_node_id "target": help_node_id
}) })

View File

@ -4,11 +4,12 @@
и переходами между узлами обработки. и переходами между узлами обработки.
""" """
from typing import Dict, Any, Optional from typing import Dict, Any, Optional, Set
from langgraph.graph import StateGraph, END from langgraph.graph import StateGraph, END
from graph_history_manager import GraphHistoryManager from graph_history_manager import GraphHistoryManager
from models import AgentState 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 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: Построение графа --- # --- LangGraph: Построение графа ---
graph_history_manager = GraphHistoryManager() graph_history_manager = GraphHistoryManager()
@ -101,6 +102,10 @@ def run_agent(user_input: str,
# ---------------------------------------------------- Initial State Setup ---------------------------------------------------------------- # ---------------------------------------------------- Initial State Setup ----------------------------------------------------------------
initial_state = AgentState(input=user_input, parent_node_id=parent_node_id) 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 графа, загружаем его историю и структуру # Если есть существующий ID графа, загружаем его историю и структуру
loaded_graph_data = None loaded_graph_data = None
if existing_graph_id: if existing_graph_id:
@ -115,6 +120,10 @@ def run_agent(user_input: str,
initial_state.graph_nodes = loaded_graph_data.get("graph_nodes", []) initial_state.graph_nodes = loaded_graph_data.get("graph_nodes", [])
initial_state.graph_edges = loaded_graph_data.get("graph_edges", []) 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 для нового входного узла
# Если parent_node_id был передан в запросе, используем его, # Если parent_node_id был передан в запросе, используем его,
# иначе берем current_node_id из загруженного графа. # иначе берем current_node_id из загруженного графа.
@ -123,6 +132,10 @@ def run_agent(user_input: str,
# История чата для текущего раунда работы LLM формируется из загруженных сообщений. # История чата для текущего раунда работы LLM формируется из загруженных сообщений.
initial_state.temporary_chat_history = loaded_graph_data.get("messages", []) initial_state.temporary_chat_history = loaded_graph_data.get("messages", [])
# Определяем graph_id, который будет использоваться для сохранения.
# Если существующий ID не передан, GraphHistoryManager сгенерирует новый.
graph_id_to_save = existing_graph_id
# ---------------------------------------------------- run_agent ---------------------------------------------------------------- # ---------------------------------------------------- run_agent ----------------------------------------------------------------
# Запускаем граф # Запускаем граф
@ -130,23 +143,25 @@ def run_agent(user_input: str,
final_state = AgentState(**result) final_state = AgentState(**result)
# ---------------------------------------------------- Save Graph ---------------------------------------------------------------- # ---------------------------------------------------- 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 больше не сохраняется напрямую, # Обратите внимание: chat_history больше не сохраняется напрямую,
# она является частью данных узлов в graph_nodes. # она является частью данных узлов в graph_nodes.
graph_data_to_save = { final_graph_id = graph_history_manager.save_graph_changes(
"id": existing_graph_id, # Используем существующий ID или он будет сгенерирован graph_id_to_save,
# менеджером истории при первом сохранении newly_added_nodes,
"graph_nodes": final_state.graph_nodes, newly_added_edges,
"graph_edges": final_state.graph_edges, final_state.current_node_id
"current_node_id": final_state.current_node_id )
}
new_graph_id = graph_history_manager.save_graph(graph_data_to_save)
# ---------------------------------------------------- prepare response ---------------------------------------------------------------- # ---------------------------------------------------- prepare response ----------------------------------------------------------------
# Для `messages` в `response_data` используем `final_state.temporary_chat_history`, # Для `messages` в `response_data` используем `final_state.temporary_chat_history`,
# так как она отражает только сообщения текущей ветки, добавленные в ходе этого выполнения. # так как она отражает только сообщения текущей ветки, добавленные в ходе этого выполнения.
response_data = { response_data = {
"graph_id": new_graph_id, "graph_id": final_graph_id,
"messages": final_state.temporary_chat_history, "messages": final_state.temporary_chat_history,
"llm_response": final_state.llm_response, "llm_response": final_state.llm_response,
"image_urls": final_state.image_urls, "image_urls": final_state.image_urls,

Binary file not shown.