/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/primitives/memory.py
831 строка
28 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""Memory primitives — базовые типы и интерфейсы для работы с памятью агентов. Использует: - ``src.primitives._time``: timezone-aware UTC, безопасный парсинг дат - ``StrEnum`` (Python 3.11+) - Pydantic ``Field`` для ``MemoryContext`` (наследует ``RunContext``) """ from __future__ import annotations from dataclasses import dataclass, field from datetime import datetime from enum import StrEnum from typing import Any, Protocol from uuid import uuid4 from pydantic import Field from src.primitives._time import ( ensure_aware, parse_datetime, to_iso, utc_now, ) from src.primitives._time import ( is_expired as check_expiry, ) from src.primitives.context import RunContext # ============================================================================ # Memory Types # ============================================================================ class MemoryType(StrEnum): """Типы памяти.""" SHORT_TERM = "short_term" # Кратковременная (в рамках сессии) LONG_TERM = "long_term" # Долговременная (персистентная) EPISODIC = "episodic" # Эпизодическая (события, диалоги) SEMANTIC = "semantic" # Семантическая (факты, знания) PROCEDURAL = "procedural" # Процедурная (навыки, how-to) USER_PROFILE = "user_profile" # Профиль пользователя class MemoryScope(StrEnum): """Область видимости памяти.""" SESSION = "session" # Только в рамках сессии USER = "user" # Для конкретного пользователя WORKSPACE = "workspace" # Для всего workspace GLOBAL = "global" # Глобальная (shared) class MemoryImportance(StrEnum): """Важность записи в памяти.""" LOW = "low" MEDIUM = "medium" HIGH = "high" CRITICAL = "critical" # ============================================================================ # Scoring / Sorting constants # ============================================================================ # Явный порядок важности (для корректной сортировки — не лексикографической!) _IMPORTANCE_ORDER: dict[MemoryImportance, int] = { MemoryImportance.LOW: 1, MemoryImportance.MEDIUM: 2, MemoryImportance.HIGH: 3, MemoryImportance.CRITICAL: 4, } # Нормализованные веса важности (для scorer, 0.0–1.0) _IMPORTANCE_WEIGHTS: dict[MemoryImportance, float] = { MemoryImportance.LOW: 0.25, MemoryImportance.MEDIUM: 0.5, MemoryImportance.HIGH: 0.75, MemoryImportance.CRITICAL: 1.0, } # Веса компонентов релевантности (сумма = 1.0) _WEIGHT_RELEVANCE = 0.4 _WEIGHT_FRESHNESS = 0.2 _WEIGHT_IMPORTANCE = 0.2 _WEIGHT_ACCESS = 0.2 # Параметры нормализации _FRESHNESS_DECAY_DAYS = 365.0 _ACCESS_NORMALIZER = 100.0 # ============================================================================ # Memory Entry # ============================================================================ @dataclass class MemoryEntry: """ Запись в памяти. Представляет единицу информации, сохранённую в памяти. """ id: str = field(default_factory=lambda: str(uuid4())) # Содержание content: str = "" content_type: str = "text" # text, json, code, etc. # Метаданные memory_type: MemoryType = MemoryType.SHORT_TERM scope: MemoryScope = MemoryScope.SESSION importance: MemoryImportance = MemoryImportance.MEDIUM # Контекст workspace_id: str = "default" user_id: str | None = None session_id: str | None = None agent_id: str | None = None chat_id: str | None = None # Теги и категории tags: list[str] = field(default_factory=list) category: str | None = None # Временные метки created_at: datetime = field(default_factory=utc_now) updated_at: datetime = field(default_factory=utc_now) expires_at: datetime | None = None # Метрики access_count: int = 0 last_accessed_at: datetime | None = None relevance_score: float = 1.0 # Дополнительные данные metadata: dict[str, Any] = field(default_factory=dict) embedding: list[float] | None = None # Векторное представление def is_expired(self) -> bool: """Проверить, истёк ли срок действия.""" return check_expiry(self.expires_at) def increment_access(self) -> None: """Инкрементировать счётчик обращений.""" self.access_count += 1 self.last_accessed_at = utc_now() def update_relevance(self, score: float) -> None: """Обновить оценку релевантности (clamp 0.0–1.0).""" self.relevance_score = max(0.0, min(1.0, score)) def add_tag(self, tag: str) -> None: """Добавить тег.""" if tag not in self.tags: self.tags.append(tag) def remove_tag(self, tag: str) -> None: """Удалить тег.""" if tag in self.tags: self.tags.remove(tag) def has_tag(self, tag: str) -> bool: """Проверить наличие тега.""" return tag in self.tags def get_metadata(self, key: str, default: Any = None) -> Any: """Получить значение из metadata.""" return self.metadata.get(key, default) def set_metadata(self, key: str, value: Any) -> None: """Установить значение в metadata.""" self.metadata[key] = value def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "content": self.content, "content_type": self.content_type, "memory_type": self.memory_type.value, "scope": self.scope.value, "importance": self.importance.value, "workspace_id": self.workspace_id, "user_id": self.user_id, "session_id": self.session_id, "agent_id": self.agent_id, "chat_id": self.chat_id, "tags": self.tags, "category": self.category, "created_at": to_iso(self.created_at), "updated_at": to_iso(self.updated_at), "expires_at": to_iso(self.expires_at), "access_count": self.access_count, "last_accessed_at": to_iso(self.last_accessed_at), "relevance_score": self.relevance_score, "metadata": self.metadata, "embedding": self.embedding, } @classmethod def from_dict(cls, data: dict[str, Any]) -> MemoryEntry: """Создать из словаря.""" return cls( id=data.get("id", str(uuid4())), content=data.get("content", ""), content_type=data.get("content_type", "text"), memory_type=MemoryType(data.get("memory_type", "short_term")), scope=MemoryScope(data.get("scope", "session")), importance=MemoryImportance(data.get("importance", "medium")), workspace_id=data.get("workspace_id", "default"), user_id=data.get("user_id"), session_id=data.get("session_id"), agent_id=data.get("agent_id"), chat_id=data.get("chat_id"), tags=data.get("tags", []), category=data.get("category"), created_at=parse_datetime(data.get("created_at")) or utc_now(), updated_at=parse_datetime(data.get("updated_at")) or utc_now(), expires_at=parse_datetime(data.get("expires_at")), access_count=data.get("access_count", 0), last_accessed_at=parse_datetime(data.get("last_accessed_at")), relevance_score=data.get("relevance_score", 1.0), metadata=data.get("metadata", {}), embedding=data.get("embedding"), ) # ============================================================================ # Memory Query # ============================================================================ @dataclass class MemoryQuery: """ Запрос к памяти. Используется для поиска записей в памяти. """ # Поисковый запрос query: str = "" query_type: str = "text" # text, semantic, keyword # Фильтры memory_types: list[MemoryType] | None = None scopes: list[MemoryScope] | None = None importance_levels: list[MemoryImportance] | None = None workspace_id: str | None = None user_id: str | None = None session_id: str | None = None agent_id: str | None = None chat_id: str | None = None tags: list[str] | None = None category: str | None = None # Временные фильтры created_after: datetime | None = None created_before: datetime | None = None # Параметры поиска limit: int = 10 offset: int = 0 min_relevance: float = 0.0 # Сортировка sort_by: str = "relevance" # relevance, created_at, updated_at, access_count sort_order: str = "desc" # asc, desc # Векторный поиск embedding: list[float] | None = None top_k: int = 10 def add_filter_type(self, memory_type: MemoryType) -> None: """Добавить фильтр по типу памяти.""" if self.memory_types is None: self.memory_types = [] if memory_type not in self.memory_types: self.memory_types.append(memory_type) def add_filter_scope(self, scope: MemoryScope) -> None: """Добавить фильтр по области видимости.""" if self.scopes is None: self.scopes = [] if scope not in self.scopes: self.scopes.append(scope) def add_filter_tag(self, tag: str) -> None: """Добавить фильтр по тегу.""" if self.tags is None: self.tags = [] if tag not in self.tags: self.tags.append(tag) def matches_entry(self, entry: MemoryEntry) -> bool: """Проверить, соответствует ли запись фильтрам.""" if self.memory_types and entry.memory_type not in self.memory_types: return False if self.scopes and entry.scope not in self.scopes: return False if self.importance_levels and entry.importance not in self.importance_levels: return False if self.workspace_id and entry.workspace_id != self.workspace_id: return False if self.user_id and entry.user_id != self.user_id: return False if self.session_id and entry.session_id != self.session_id: return False if self.agent_id and entry.agent_id != self.agent_id: return False if self.chat_id and entry.chat_id != self.chat_id: return False if self.tags and not any(entry.has_tag(tag) for tag in self.tags): return False if self.category and entry.category != self.category: return False # Временные фильтры (timezone-aware сравнение) if self.created_after and ensure_aware(entry.created_at) < ensure_aware(self.created_after): return False if self.created_before and ensure_aware(entry.created_at) > ensure_aware( self.created_before ): return False if entry.relevance_score < self.min_relevance: return False if entry.is_expired(): return False return True def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "query": self.query, "query_type": self.query_type, "memory_types": [mt.value for mt in self.memory_types] if self.memory_types else None, "scopes": [s.value for s in self.scopes] if self.scopes else None, "importance_levels": [il.value for il in self.importance_levels] if self.importance_levels else None, "workspace_id": self.workspace_id, "user_id": self.user_id, "session_id": self.session_id, "agent_id": self.agent_id, "chat_id": self.chat_id, "tags": self.tags, "category": self.category, "created_after": to_iso(self.created_after), "created_before": to_iso(self.created_before), "limit": self.limit, "offset": self.offset, "min_relevance": self.min_relevance, "sort_by": self.sort_by, "sort_order": self.sort_order, "embedding": self.embedding, "top_k": self.top_k, } # ============================================================================ # Memory Result # ============================================================================ @dataclass class MemoryResult: """ Результат поиска в памяти. Содержит найденные записи и метаданные поиска. """ entries: list[MemoryEntry] = field(default_factory=list) total_count: int = 0 query_time_ms: float = 0.0 # Метаданные поиска query: MemoryQuery | None = None def add_entry(self, entry: MemoryEntry) -> None: """Добавить запись.""" self.entries.append(entry) self.total_count = len(self.entries) def get_top_entries(self, n: int) -> list[MemoryEntry]: """Получить top N записей.""" return self.entries[:n] def is_empty(self) -> bool: """Проверить, пуст ли результат.""" return len(self.entries) == 0 def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "entries": [entry.to_dict() for entry in self.entries], "total_count": self.total_count, "query_time_ms": self.query_time_ms, "query": self.query.to_dict() if self.query else None, } # ============================================================================ # Memory Context # ============================================================================ class MemoryContext(RunContext): """ Контекст для работы с памятью. Расширяет RunContext специфичными для памяти полями. Note: это Pydantic-модель (наследует RunContext) — используется Pydantic ``Field``, а не dataclass ``field``. """ memory_type: MemoryType = MemoryType.SHORT_TERM scope: MemoryScope = MemoryScope.SESSION # История операций stored_entries: list[str] = Field(default_factory=list) # IDs сохранённых записей retrieved_entries: list[str] = Field(default_factory=list) # IDs извлечённых записей # Кэш cache_hits: int = 0 cache_misses: int = 0 def add_stored_entry(self, entry_id: str) -> None: """Добавить ID сохранённой записи.""" self.stored_entries.append(entry_id) def add_retrieved_entry(self, entry_id: str) -> None: """Добавить ID извлечённой записи.""" self.retrieved_entries.append(entry_id) def increment_cache_hit(self) -> None: """Инкрементировать счётчик попаданий в кэш.""" self.cache_hits += 1 def increment_cache_miss(self) -> None: """Инкрементировать счётчик промахов кэша.""" self.cache_misses += 1 def get_cache_hit_rate(self) -> float: """Получить процент попаданий в кэш.""" total = self.cache_hits + self.cache_misses if total == 0: return 0.0 return self.cache_hits / total def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" data = super().to_dict() data.update( { "memory_type": self.memory_type.value, "scope": self.scope.value, "stored_entries": self.stored_entries, "retrieved_entries": self.retrieved_entries, "cache_hits": self.cache_hits, "cache_misses": self.cache_misses, "cache_hit_rate": self.get_cache_hit_rate(), } ) return data # ============================================================================ # Memory Store Protocol # ============================================================================ class MemoryStore(Protocol): """ Протокол для хранилища памяти. Определяет интерфейс для различных реализаций хранилищ. """ async def store(self, entry: MemoryEntry, context: MemoryContext | None = None) -> str: """Сохранить запись в памяти.""" ... async def retrieve( self, query: MemoryQuery, context: MemoryContext | None = None ) -> MemoryResult: """Извлечь записи из памяти.""" ... async def update( self, entry_id: str, updates: dict[str, Any], context: MemoryContext | None = None ) -> MemoryEntry | None: """Обновить запись в памяти.""" ... async def delete(self, entry_id: str, context: MemoryContext | None = None) -> bool: """Удалить запись из памяти.""" ... async def get(self, entry_id: str, context: MemoryContext | None = None) -> MemoryEntry | None: """Получить запись по ID.""" ... async def list( self, workspace_id: str, user_id: str | None = None, limit: int = 100, offset: int = 0, context: MemoryContext | None = None, ) -> list[MemoryEntry]: """Получить список записей.""" ... async def count( self, workspace_id: str, user_id: str | None = None, memory_type: MemoryType | None = None, context: MemoryContext | None = None, ) -> int: """Подсчитать количество записей.""" ... async def clear( self, workspace_id: str, user_id: str | None = None, memory_type: MemoryType | None = None, context: MemoryContext | None = None, ) -> int: """Очистить записи.""" ... # ============================================================================ # Memory Helpers # ============================================================================ def create_short_term_memory( content: str, workspace_id: str = "default", session_id: str | None = None, **kwargs: Any, ) -> MemoryEntry: """Создать кратковременную запись в памяти.""" return MemoryEntry( content=content, memory_type=MemoryType.SHORT_TERM, scope=MemoryScope.SESSION, workspace_id=workspace_id, session_id=session_id, **kwargs, ) def create_long_term_memory( content: str, workspace_id: str = "default", user_id: str | None = None, **kwargs: Any, ) -> MemoryEntry: """Создать долговременную запись в памяти.""" return MemoryEntry( content=content, memory_type=MemoryType.LONG_TERM, scope=MemoryScope.USER, workspace_id=workspace_id, user_id=user_id, **kwargs, ) def create_episodic_memory( content: str, workspace_id: str = "default", user_id: str | None = None, chat_id: str | None = None, **kwargs: Any, ) -> MemoryEntry: """Создать эпизодическую запись в памяти.""" return MemoryEntry( content=content, memory_type=MemoryType.EPISODIC, scope=MemoryScope.USER, workspace_id=workspace_id, user_id=user_id, chat_id=chat_id, **kwargs, ) def create_semantic_memory( content: str, workspace_id: str = "default", **kwargs: Any, ) -> MemoryEntry: """Создать семантическую запись в памяти.""" return MemoryEntry( content=content, memory_type=MemoryType.SEMANTIC, scope=MemoryScope.WORKSPACE, workspace_id=workspace_id, **kwargs, ) def create_procedural_memory( content: str, workspace_id: str = "default", **kwargs: Any, ) -> MemoryEntry: """Создать процедурную запись в памяти.""" return MemoryEntry( content=content, memory_type=MemoryType.PROCEDURAL, scope=MemoryScope.WORKSPACE, workspace_id=workspace_id, **kwargs, ) def create_user_profile_memory( content: str, workspace_id: str = "default", user_id: str | None = None, **kwargs: Any, ) -> MemoryEntry: """Создать запись профиля пользователя.""" return MemoryEntry( content=content, memory_type=MemoryType.USER_PROFILE, scope=MemoryScope.USER, workspace_id=workspace_id, user_id=user_id, **kwargs, ) # ============================================================================ # Memory Scorer # ============================================================================ class MemoryScorer: """ Оценка релевантности записей в памяти. Используется для ранжирования результатов поиска. Веса компонентов нормализованы (сумма = 1.0). """ @staticmethod def calculate_relevance( entry: MemoryEntry, query: MemoryQuery, current_time: datetime | None = None, ) -> float: """ Вычислить оценку релевантности записи. Компоненты (веса нормализованы, сумма = 1.0): - Базовая релевантность (0.4) - Свежесть (0.2) — убывает в течение года - Важность (0.2) — нормализована 0.25–1.0 - Частота доступа (0.2) Returns: Оценка релевантности от 0.0 до 1.0. """ if current_time is None: current_time = utc_now() score = 0.0 # Базовая релевантность (от embedding similarity или keyword match) score += entry.relevance_score * _WEIGHT_RELEVANCE # Свежесть (более новые записи получают больший вес) age_days = (ensure_aware(current_time) - ensure_aware(entry.created_at)).days freshness_score = max(0.0, 1.0 - (age_days / _FRESHNESS_DECAY_DAYS)) score += freshness_score * _WEIGHT_FRESHNESS # Важность (нормализовано) score += _IMPORTANCE_WEIGHTS.get(entry.importance, 0.5) * _WEIGHT_IMPORTANCE # Частота доступа (популярные записи получают больший вес) access_score = min(1.0, entry.access_count / _ACCESS_NORMALIZER) score += access_score * _WEIGHT_ACCESS return max(0.0, min(1.0, score)) @staticmethod def rank_entries( entries: list[MemoryEntry], query: MemoryQuery, current_time: datetime | None = None, ) -> list[MemoryEntry]: """ Ранжировать записи по релевантности. Returns: Отсортированный список записей (по убыванию релевантности). """ scored_entries = [ (entry, MemoryScorer.calculate_relevance(entry, query, current_time)) for entry in entries ] scored_entries.sort(key=lambda x: x[1], reverse=True) # Обновление relevance_score в записях for entry, score in scored_entries: entry.update_relevance(score) return [entry for entry, _ in scored_entries] # ============================================================================ # Memory Summarizer # ============================================================================ class MemorySummarizer: """ Суммаризация записей в памяти. Используется для создания кратких обзоров из множества записей. """ @staticmethod def summarize(entries: list[MemoryEntry], max_length: int = 500) -> str: """ Создать краткое резюме из записей. Сортировка по важности (явный порядок, НЕ лексикографический) и релевантности. Returns: Строка с резюме. """ if not entries: return "" # ✅ Корректная сортировка: явный порядок важности + релевантность sorted_entries = sorted( entries, key=lambda e: (_IMPORTANCE_ORDER[e.importance], e.relevance_score), reverse=True, ) summaries: list[str] = [] current_length = 0 for entry in sorted_entries: content = entry.content.strip() if current_length + len(content) > max_length: break summaries.append(content) current_length += len(content) return "\n\n".join(summaries) @staticmethod def extract_key_facts(entries: list[MemoryEntry], max_facts: int = 5) -> list[str]: """ Извлечь ключевые факты из записей. Returns: Список ключевых фактов. """ # Фильтрация по важности important_entries = [ e for e in entries if e.importance in (MemoryImportance.HIGH, MemoryImportance.CRITICAL) ] # Если мало важных записей, берём все if len(important_entries) < max_facts: important_entries = entries # Сортировка по релевантности sorted_entries = sorted( important_entries, key=lambda e: e.relevance_score, reverse=True, ) facts: list[str] = [] for entry in sorted_entries[:max_facts]: content = entry.content.strip() first_line = content.split("\n")[0] if len(first_line) > 100: first_line = first_line[:100] + "..." facts.append(first_line) return facts # ============================================================================ # Exports # ============================================================================ __all__ = [ "MemoryContext", "MemoryEntry", "MemoryImportance", "MemoryQuery", "MemoryResult", "MemoryScope", "MemoryScorer", "MemoryStore", "MemorySummarizer", "MemoryType", "create_episodic_memory", "create_long_term_memory", "create_procedural_memory", "create_semantic_memory", "create_short_term_memory", "create_user_profile_memory", ]