/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/runtime/state.py
131 строка
5 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""Runtime State Manager — in-memory реализация StateManager protocol.""" from __future__ import annotations import asyncio from collections.abc import Callable from typing import Any import structlog from src.primitives import SessionState, StateEvent logger = structlog.get_logger() # Тип подписчика на события состояния StateSubscriber = Callable[[StateEvent], Any] class InMemoryStateManager: """ In-memory реализация StateManager (из primitives). Хранит сессии в памяти, поддерживает подписку на события. Для production замените на persistent-реализацию (Postgres/Redis). """ def __init__(self) -> None: self._sessions: dict[str, SessionState] = {} self._subscribers: list[StateSubscriber] = [] self._lock = asyncio.Lock() logger.info("state_manager.initialized", backend="in_memory") # ======================================================================== # Session CRUD # ======================================================================== async def save_session(self, session: SessionState) -> None: """Сохранить состояние сессии.""" async with self._lock: self._sessions[session.session_id] = session await self.emit_event( StateEvent( event_type="session_saved", session_id=session.session_id, new_value=session.session_id, ) ) async def load_session(self, session_id: str) -> SessionState | None: """Загрузить состояние сессии.""" async with self._lock: return self._sessions.get(session_id) async def delete_session(self, session_id: str) -> bool: """Удалить состояние сессии.""" async with self._lock: removed = self._sessions.pop(session_id, None) if removed is not None: await self.emit_event(StateEvent(event_type="session_deleted", session_id=session_id)) return True return False async def list_sessions( self, user_id: str | None = None, limit: int = 100 ) -> list[SessionState]: """Получить список сессий (опционально фильтр по user).""" async with self._lock: sessions = list(self._sessions.values()) if user_id is not None: sessions = [s for s in sessions if s.user_id == user_id] return sessions[:limit] # ======================================================================== # Events # ======================================================================== async def emit_event(self, event: StateEvent) -> None: """Отправить событие всем подписчикам.""" for callback in self._subscribers: try: result = callback(event) if asyncio.iscoroutine(result): await result except Exception as e: logger.warning("state_manager.subscriber_error", error=str(e)) async def subscribe(self, callback: StateSubscriber) -> None: """Подписаться на события состояния.""" if callback not in self._subscribers: self._subscribers.append(callback) async def unsubscribe(self, callback: StateSubscriber) -> None: """Отписаться от событий.""" if callback in self._subscribers: self._subscribers.remove(callback) # ======================================================================== # Utilities # ======================================================================== async def get_or_create_session( self, session_id: str, user_id: str | None = None, workspace_id: str = "default" ) -> SessionState: """Получить сессию или создать новую.""" session = await self.load_session(session_id) if session is None: session = SessionState( session_id=session_id, user_id=user_id, workspace_id=workspace_id ) await self.save_session(session) return session async def cleanup_inactive(self, timeout_minutes: int = 30) -> int: """Очистить неактивные сессии. Возвращает количество удалённых.""" async with self._lock: inactive = [ sid for sid, s in self._sessions.items() if not s.is_active(timeout_minutes) ] for sid in inactive: del self._sessions[sid] if inactive: logger.info("state_manager.cleanup", removed=len(inactive)) return len(inactive) async def count(self) -> int: """Количество активных сессий.""" async with self._lock: return len(self._sessions) __all__ = ["InMemoryStateManager", "StateSubscriber"]