/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/primitives/state.py
1 072 строки
38 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""Состояние агентов, flows и сессий. Использует: - ``src.primitives._time``: timezone-aware UTC, monotonic для duration - ``StrEnum`` (Python 3.11+) для статусов - ``copy.deepcopy`` для checkpoint (корректное копирование вложенных данных) """ from __future__ import annotations import copy from collections.abc import Callable from dataclasses import dataclass, field from datetime import datetime from enum import StrEnum from typing import Any, Protocol from uuid import uuid4 from src.primitives._time import ( elapsed_ms, ensure_aware, monotonic, parse_datetime, to_iso, utc_now, ) # ============================================================================ # Status Enums # ============================================================================ class TaskStatus(StrEnum): """Статус выполнения задачи.""" PENDING = "pending" RUNNING = "running" PAUSED = "paused" COMPLETED = "completed" FAILED = "failed" CANCELLED = "cancelled" TIMEOUT = "timeout" class NodeStatus(StrEnum): """Статус узла в flow.""" PENDING = "pending" RUNNING = "running" COMPLETED = "completed" FAILED = "failed" SKIPPED = "skipped" WAITING = "waiting" class EdgeStatus(StrEnum): """Статус ребра в flow.""" INACTIVE = "inactive" ACTIVE = "active" TRAVERSED = "traversed" BLOCKED = "blocked" class RetryStrategy(StrEnum): """Стратегия повторных попыток.""" NONE = "none" FIXED = "fixed" EXPONENTIAL = "exponential" LINEAR = "linear" # ============================================================================ # Retry delay constants # ============================================================================ _FIXED_RETRY_DELAY_SECONDS = 1.0 _EXPONENTIAL_RETRY_BASE = 2.0 # ============================================================================ # Execution State # ============================================================================ @dataclass class ExecutionState: """ Состояние выполнения отдельного шага. Отслеживает прогресс, метрики и результаты выполнения. """ id: str = field(default_factory=lambda: str(uuid4())) step_name: str = "" status: TaskStatus = TaskStatus.PENDING # Тайминги started_at: datetime | None = None completed_at: datetime | None = None duration_ms: float = 0.0 # Метрики tokens_used: int = 0 llm_calls: int = 0 tool_calls: int = 0 retry_count: int = 0 # Данные input_data: dict[str, Any] = field(default_factory=dict) output_data: dict[str, Any] = field(default_factory=dict) error: str | None = None # Retry retry_strategy: RetryStrategy = RetryStrategy.NONE max_retries: int = 3 _start_time: float = field(default_factory=monotonic, repr=False) def start(self) -> None: """Начать выполнение.""" self.status = TaskStatus.RUNNING self.started_at = utc_now() self._start_time = monotonic() def complete(self, output: dict[str, Any] | None = None) -> None: """Завершить успешно.""" self.status = TaskStatus.COMPLETED self.completed_at = utc_now() self.duration_ms = elapsed_ms(self._start_time) if output: self.output_data = output def fail(self, error: str) -> None: """Завершить с ошибкой.""" self.status = TaskStatus.FAILED self.completed_at = utc_now() self.duration_ms = elapsed_ms(self._start_time) self.error = error def cancel(self) -> None: """Отменить.""" self.status = TaskStatus.CANCELLED self.completed_at = utc_now() self.duration_ms = elapsed_ms(self._start_time) def pause(self) -> None: """Приостановить.""" self.status = TaskStatus.PAUSED self.duration_ms = elapsed_ms(self._start_time) def timeout(self) -> None: """Таймаут.""" self.status = TaskStatus.TIMEOUT self.completed_at = utc_now() self.duration_ms = elapsed_ms(self._start_time) def increment_retry(self) -> None: """Инкрементировать счётчик попыток.""" self.retry_count += 1 def can_retry(self) -> bool: """Проверить, можно ли повторить.""" if self.retry_strategy == RetryStrategy.NONE: return False return self.retry_count < self.max_retries def get_retry_delay(self) -> float: """Получить задержку перед повтором в секундах.""" if self.retry_strategy == RetryStrategy.FIXED: return _FIXED_RETRY_DELAY_SECONDS if self.retry_strategy == RetryStrategy.LINEAR: return float(self.retry_count + 1) if self.retry_strategy == RetryStrategy.EXPONENTIAL: return _EXPONENTIAL_RETRY_BASE**self.retry_count return 0.0 def add_tokens(self, count: int) -> None: """Добавить использованные токены.""" self.tokens_used += count def add_llm_call(self) -> None: """Добавить LLM вызов.""" self.llm_calls += 1 def add_tool_call(self) -> None: """Добавить tool вызов.""" self.tool_calls += 1 def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "step_name": self.step_name, "status": self.status.value, "started_at": to_iso(self.started_at), "completed_at": to_iso(self.completed_at), "duration_ms": self.duration_ms, "tokens_used": self.tokens_used, "llm_calls": self.llm_calls, "tool_calls": self.tool_calls, "retry_count": self.retry_count, "input_data": self.input_data, "output_data": self.output_data, "error": self.error, "retry_strategy": self.retry_strategy.value, "max_retries": self.max_retries, } @classmethod def from_dict(cls, data: dict[str, Any]) -> ExecutionState: """Создать из словаря.""" return cls( id=data.get("id", str(uuid4())), step_name=data.get("step_name", ""), status=TaskStatus(data.get("status", "pending")), started_at=parse_datetime(data.get("started_at")), completed_at=parse_datetime(data.get("completed_at")), duration_ms=data.get("duration_ms", 0.0), tokens_used=data.get("tokens_used", 0), llm_calls=data.get("llm_calls", 0), tool_calls=data.get("tool_calls", 0), retry_count=data.get("retry_count", 0), input_data=data.get("input_data", {}), output_data=data.get("output_data", {}), error=data.get("error"), retry_strategy=RetryStrategy(data.get("retry_strategy", "none")), max_retries=data.get("max_retries", 3), ) # ============================================================================ # Agent State # ============================================================================ @dataclass class AgentState: """ Состояние агента во время выполнения. Содержит полную информацию о текущем состоянии агента: идентификация, история сообщений, метрики, инструменты, память. """ agent_id: str agent_name: str status: TaskStatus = TaskStatus.PENDING # Данные input_data: dict[str, Any] = field(default_factory=dict) output_data: dict[str, Any] = field(default_factory=dict) messages: list[dict[str, Any]] = field(default_factory=list) # Метрики tokens_used: int = 0 duration_ms: float = 0.0 llm_calls: int = 0 tool_calls: int = 0 # Ошибки error: str | None = None error_count: int = 0 # Инструменты available_tools: list[str] = field(default_factory=list) used_tools: list[str] = field(default_factory=list) # Память short_term_memory: list[str] = field(default_factory=list) long_term_memory: dict[str, Any] = field(default_factory=dict) # Конфигурация system_prompt: str | None = None model: str = "deepseek-ai/DeepSeek-V4-Pro" temperature: float = 0.7 max_tokens: int = 4096 # Тайминги started_at: datetime | None = None completed_at: datetime | None = None _start_time: float = field(default_factory=monotonic, repr=False) def start(self) -> None: """Начать выполнение.""" self.status = TaskStatus.RUNNING self.started_at = utc_now() self._start_time = monotonic() def complete(self, output: dict[str, Any] | None = None) -> None: """Завершить успешно.""" self.status = TaskStatus.COMPLETED self.completed_at = utc_now() self.duration_ms = elapsed_ms(self._start_time) if output: self.output_data = output def fail(self, error: str) -> None: """Завершить с ошибкой.""" self.status = TaskStatus.FAILED self.completed_at = utc_now() self.duration_ms = elapsed_ms(self._start_time) self.error = error self.error_count += 1 def cancel(self) -> None: """Отменить.""" self.status = TaskStatus.CANCELLED self.completed_at = utc_now() self.duration_ms = elapsed_ms(self._start_time) def add_message(self, role: str, content: str, **kwargs: Any) -> None: """Добавить сообщение в историю.""" self.messages.append( { "role": role, "content": content, "timestamp": utc_now().isoformat(), **kwargs, } ) def get_messages(self, role: str | None = None) -> list[dict[str, Any]]: """Получить сообщения, опционально фильтруя по роли.""" if role is None: return self.messages return [m for m in self.messages if m.get("role") == role] def get_last_message(self, role: str | None = None) -> dict[str, Any] | None: """Получить последнее сообщение.""" messages = self.get_messages(role) return messages[-1] if messages else None def clear_messages(self) -> None: """Очистить историю сообщений.""" self.messages.clear() def set_output(self, key: str, value: Any) -> None: """Установить выходное значение.""" self.output_data[key] = value def get_output(self, key: str, default: Any = None) -> Any: """Получить выходное значение.""" return self.output_data.get(key, default) def add_tokens(self, count: int) -> None: """Добавить использованные токены.""" self.tokens_used += count def add_llm_call(self) -> None: """Добавить LLM вызов.""" self.llm_calls += 1 def add_tool_call(self, tool_name: str) -> None: """Добавить tool вызов.""" self.tool_calls += 1 if tool_name not in self.used_tools: self.used_tools.append(tool_name) def add_tool(self, tool_name: str) -> None: """Добавить доступный инструмент.""" if tool_name not in self.available_tools: self.available_tools.append(tool_name) def remove_tool(self, tool_name: str) -> None: """Удалить инструмент.""" if tool_name in self.available_tools: self.available_tools.remove(tool_name) def has_tool(self, tool_name: str) -> bool: """Проверить наличие инструмента.""" return tool_name in self.available_tools def add_to_short_term_memory(self, content: str) -> None: """Добавить в кратковременную память (с ограничением размера).""" self.short_term_memory.append(content) if len(self.short_term_memory) > 100: self.short_term_memory = self.short_term_memory[-100:] def set_long_term_memory(self, key: str, value: Any) -> None: """Установить долговременную память.""" self.long_term_memory[key] = value def get_long_term_memory(self, key: str, default: Any = None) -> Any: """Получить из долговременной памяти.""" return self.long_term_memory.get(key, default) def clear_short_term_memory(self) -> None: """Очистить кратковременную память.""" self.short_term_memory.clear() def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "agent_id": self.agent_id, "agent_name": self.agent_name, "status": self.status.value, "input_data": self.input_data, "output_data": self.output_data, "messages": self.messages, "tokens_used": self.tokens_used, "duration_ms": self.duration_ms, "llm_calls": self.llm_calls, "tool_calls": self.tool_calls, "error": self.error, "error_count": self.error_count, "available_tools": self.available_tools, "used_tools": self.used_tools, "short_term_memory": self.short_term_memory, "long_term_memory": self.long_term_memory, "system_prompt": self.system_prompt, "model": self.model, "temperature": self.temperature, "max_tokens": self.max_tokens, "started_at": to_iso(self.started_at), "completed_at": to_iso(self.completed_at), } @classmethod def from_dict(cls, data: dict[str, Any]) -> AgentState: """Создать из словаря.""" return cls( agent_id=data["agent_id"], agent_name=data["agent_name"], status=TaskStatus(data.get("status", "pending")), input_data=data.get("input_data", {}), output_data=data.get("output_data", {}), messages=data.get("messages", []), tokens_used=data.get("tokens_used", 0), duration_ms=data.get("duration_ms", 0.0), llm_calls=data.get("llm_calls", 0), tool_calls=data.get("tool_calls", 0), error=data.get("error"), error_count=data.get("error_count", 0), available_tools=data.get("available_tools", []), used_tools=data.get("used_tools", []), short_term_memory=data.get("short_term_memory", []), long_term_memory=data.get("long_term_memory", {}), system_prompt=data.get("system_prompt"), model=data.get("model", "deepseek-ai/DeepSeek-V4-Pro"), temperature=data.get("temperature", 0.7), max_tokens=data.get("max_tokens", 4096), started_at=parse_datetime(data.get("started_at")), completed_at=parse_datetime(data.get("completed_at")), ) # ============================================================================ # Checkpoint # ============================================================================ @dataclass class Checkpoint: """ Checkpoint для сохранения состояния flow. Используется для возможности отката к предыдущему состоянию. """ id: str = field(default_factory=lambda: str(uuid4())) flow_id: str = "" name: str = "" # Сохранённое состояние agent_states: dict[str, AgentState] = field(default_factory=dict) shared_data: dict[str, Any] = field(default_factory=dict) variables: dict[str, Any] = field(default_factory=dict) # Метрики total_tokens: int = 0 total_duration_ms: float = 0.0 # Metadata created_at: datetime = field(default_factory=utc_now) description: str = "" def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "flow_id": self.flow_id, "name": self.name, "agent_states": {k: v.to_dict() for k, v in self.agent_states.items()}, "shared_data": self.shared_data, "variables": self.variables, "total_tokens": self.total_tokens, "total_duration_ms": self.total_duration_ms, "created_at": to_iso(self.created_at), "description": self.description, } @classmethod def from_dict(cls, data: dict[str, Any]) -> Checkpoint: """Создать из словаря.""" return cls( id=data.get("id", str(uuid4())), flow_id=data.get("flow_id", ""), name=data.get("name", ""), agent_states={ k: AgentState.from_dict(v) for k, v in data.get("agent_states", {}).items() }, shared_data=data.get("shared_data", {}), variables=data.get("variables", {}), total_tokens=data.get("total_tokens", 0), total_duration_ms=data.get("total_duration_ms", 0.0), created_at=parse_datetime(data.get("created_at")) or utc_now(), description=data.get("description", ""), ) # ============================================================================ # Flow State # ============================================================================ @dataclass class FlowState: """ Состояние всего flow. Содержит информацию о выполнении всего flow: состояния агентов, общие данные, переменные, checkpoints, метрики. """ flow_id: str flow_name: str status: TaskStatus = TaskStatus.PENDING # Состояния агентов agent_states: dict[str, AgentState] = field(default_factory=dict) # Общие данные shared_data: dict[str, Any] = field(default_factory=dict) variables: dict[str, Any] = field(default_factory=dict) # Метрики total_tokens: int = 0 total_duration_ms: float = 0.0 total_llm_calls: int = 0 total_tool_calls: int = 0 # Ошибки error: str | None = None error_count: int = 0 # Execution history execution_history: list[ExecutionState] = field(default_factory=list) # Checkpoints checkpoints: list[Checkpoint] = field(default_factory=list) # Тайминги started_at: datetime | None = None completed_at: datetime | None = None _start_time: float = field(default_factory=monotonic, repr=False) def start(self) -> None: """Начать выполнение.""" self.status = TaskStatus.RUNNING self.started_at = utc_now() self._start_time = monotonic() def complete(self) -> None: """Завершить успешно.""" self.status = TaskStatus.COMPLETED self.completed_at = utc_now() self.total_duration_ms = elapsed_ms(self._start_time) self._aggregate_metrics() def fail(self, error: str) -> None: """Завершить с ошибкой.""" self.status = TaskStatus.FAILED self.completed_at = utc_now() self.total_duration_ms = elapsed_ms(self._start_time) self.error = error self.error_count += 1 self._aggregate_metrics() def cancel(self) -> None: """Отменить.""" self.status = TaskStatus.CANCELLED self.completed_at = utc_now() self.total_duration_ms = elapsed_ms(self._start_time) self._aggregate_metrics() def _aggregate_metrics(self) -> None: """Агрегировать метрики из всех агентов.""" self.total_tokens = sum(agent.tokens_used for agent in self.agent_states.values()) self.total_llm_calls = sum(agent.llm_calls for agent in self.agent_states.values()) self.total_tool_calls = sum(agent.tool_calls for agent in self.agent_states.values()) def add_agent_state(self, agent_state: AgentState) -> None: """Добавить состояние агента.""" self.agent_states[agent_state.agent_id] = agent_state def get_agent_state(self, agent_id: str) -> AgentState | None: """Получить состояние агента.""" return self.agent_states.get(agent_id) def remove_agent_state(self, agent_id: str) -> bool: """Удалить состояние агента.""" if agent_id in self.agent_states: del self.agent_states[agent_id] return True return False def set_shared(self, key: str, value: Any) -> None: """Установить общие данные.""" self.shared_data[key] = value def get_shared(self, key: str, default: Any = None) -> Any: """Получить общие данные.""" return self.shared_data.get(key, default) def delete_shared(self, key: str) -> None: """Удалить общие данные.""" if key in self.shared_data: del self.shared_data[key] def set_variable(self, name: str, value: Any) -> None: """Установить переменную flow.""" self.variables[name] = value def get_variable(self, name: str, default: Any = None) -> Any: """Получить переменную flow.""" return self.variables.get(name, default) def delete_variable(self, name: str) -> None: """Удалить переменную.""" if name in self.variables: del self.variables[name] def has_variable(self, name: str) -> bool: """Проверить наличие переменной.""" return name in self.variables def add_execution_step(self, step: ExecutionState) -> None: """Добавить шаг выполнения.""" self.execution_history.append(step) def get_last_execution_step(self) -> ExecutionState | None: """Получить последний шаг выполнения.""" return self.execution_history[-1] if self.execution_history else None def get_execution_steps_by_status(self, status: TaskStatus) -> list[ExecutionState]: """Получить шаги выполнения по статусу.""" return [s for s in self.execution_history if s.status == status] def create_checkpoint(self, name: str = "") -> Checkpoint: """ Создать checkpoint для возможности отката. Использует ``copy.deepcopy`` для корректного копирования вложенных структур (shared_data, variables) — изменения в оригинале не влияют на checkpoint и наоборот. """ checkpoint = Checkpoint( flow_id=self.flow_id, name=name, agent_states=copy.deepcopy(self.agent_states), shared_data=copy.deepcopy(self.shared_data), variables=copy.deepcopy(self.variables), total_tokens=self.total_tokens, total_duration_ms=self.total_duration_ms, ) self.checkpoints.append(checkpoint) return checkpoint def restore_checkpoint(self, checkpoint: Checkpoint) -> None: """ Восстановить состояние из checkpoint. Использует ``copy.deepcopy`` чтобы отвязать состояние от checkpoint. """ self.agent_states = copy.deepcopy(checkpoint.agent_states) self.shared_data = copy.deepcopy(checkpoint.shared_data) self.variables = copy.deepcopy(checkpoint.variables) self.total_tokens = checkpoint.total_tokens self.total_duration_ms = checkpoint.total_duration_ms def get_checkpoint(self, name: str) -> Checkpoint | None: """Получить checkpoint по имени.""" return next((c for c in self.checkpoints if c.name == name), None) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "flow_id": self.flow_id, "flow_name": self.flow_name, "status": self.status.value, "agent_states": {k: v.to_dict() for k, v in self.agent_states.items()}, "shared_data": self.shared_data, "variables": self.variables, "total_tokens": self.total_tokens, "total_duration_ms": self.total_duration_ms, "total_llm_calls": self.total_llm_calls, "total_tool_calls": self.total_tool_calls, "error": self.error, "error_count": self.error_count, "execution_history": [s.to_dict() for s in self.execution_history], "checkpoints": [c.to_dict() for c in self.checkpoints], "started_at": to_iso(self.started_at), "completed_at": to_iso(self.completed_at), } @classmethod def from_dict(cls, data: dict[str, Any]) -> FlowState: """Создать из словаря.""" return cls( flow_id=data["flow_id"], flow_name=data["flow_name"], status=TaskStatus(data.get("status", "pending")), agent_states={ k: AgentState.from_dict(v) for k, v in data.get("agent_states", {}).items() }, shared_data=data.get("shared_data", {}), variables=data.get("variables", {}), total_tokens=data.get("total_tokens", 0), total_duration_ms=data.get("total_duration_ms", 0.0), total_llm_calls=data.get("total_llm_calls", 0), total_tool_calls=data.get("total_tool_calls", 0), error=data.get("error"), error_count=data.get("error_count", 0), execution_history=[ ExecutionState.from_dict(s) for s in data.get("execution_history", []) ], checkpoints=[Checkpoint.from_dict(c) for c in data.get("checkpoints", [])], started_at=parse_datetime(data.get("started_at")), completed_at=parse_datetime(data.get("completed_at")), ) # ============================================================================ # Session State # ============================================================================ @dataclass class SessionState: """ Состояние пользовательской сессии. Управляет состоянием в рамках одной сессии пользователя: активные flows и агенты, история взаимодействий, настройки, контекст. """ session_id: str user_id: str | None = None workspace_id: str = "default" # Активные выполнения active_flows: dict[str, FlowState] = field(default_factory=dict) active_agents: dict[str, AgentState] = field(default_factory=dict) # История interaction_history: list[dict[str, Any]] = field(default_factory=list) # Настройки settings: dict[str, Any] = field(default_factory=dict) # Контекст context: dict[str, Any] = field(default_factory=dict) # Статистика total_interactions: int = 0 total_flows_completed: int = 0 total_agents_completed: int = 0 # Timestamps created_at: datetime = field(default_factory=utc_now) last_activity_at: datetime = field(default_factory=utc_now) def add_active_flow(self, flow_state: FlowState) -> None: """Добавить активный flow.""" self.active_flows[flow_state.flow_id] = flow_state self.last_activity_at = utc_now() def remove_active_flow(self, flow_id: str) -> bool: """Удалить активный flow.""" if flow_id in self.active_flows: del self.active_flows[flow_id] self.last_activity_at = utc_now() return True return False def get_active_flow(self, flow_id: str) -> FlowState | None: """Получить активный flow.""" return self.active_flows.get(flow_id) def add_active_agent(self, agent_state: AgentState) -> None: """Добавить активного агента.""" self.active_agents[agent_state.agent_id] = agent_state self.last_activity_at = utc_now() def remove_active_agent(self, agent_id: str) -> bool: """Удалить активного агента.""" if agent_id in self.active_agents: del self.active_agents[agent_id] self.last_activity_at = utc_now() return True return False def get_active_agent(self, agent_id: str) -> AgentState | None: """Получить активного агента.""" return self.active_agents.get(agent_id) def add_interaction(self, interaction: dict[str, Any]) -> None: """Добавить взаимодействие в историю (с ограничением размера).""" self.interaction_history.append( { **interaction, "timestamp": utc_now().isoformat(), } ) self.total_interactions += 1 self.last_activity_at = utc_now() if len(self.interaction_history) > 1000: self.interaction_history = self.interaction_history[-1000:] def get_recent_interactions(self, limit: int = 10) -> list[dict[str, Any]]: """Получить последние взаимодействия.""" return self.interaction_history[-limit:] def set_setting(self, key: str, value: Any) -> None: """Установить настройку.""" self.settings[key] = value def get_setting(self, key: str, default: Any = None) -> Any: """Получить настройку.""" return self.settings.get(key, default) def set_context(self, key: str, value: Any) -> None: """Установить контекст.""" self.context[key] = value def get_context(self, key: str, default: Any = None) -> Any: """Получить контекст.""" return self.context.get(key, default) def is_active(self, timeout_minutes: int = 30) -> bool: """Проверить, активна ли сессия.""" diff = (utc_now() - ensure_aware(self.last_activity_at)).total_seconds() / 60 return diff < timeout_minutes def cleanup_inactive(self) -> None: """Очистить завершённые flows и агенты.""" completed_flows = [ flow_id for flow_id, flow in self.active_flows.items() if flow.status in (TaskStatus.COMPLETED, TaskStatus.FAILED, TaskStatus.CANCELLED) ] for flow_id in completed_flows: self.remove_active_flow(flow_id) completed_agents = [ agent_id for agent_id, agent in self.active_agents.items() if agent.status in (TaskStatus.COMPLETED, TaskStatus.FAILED, TaskStatus.CANCELLED) ] for agent_id in completed_agents: self.remove_active_agent(agent_id) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "session_id": self.session_id, "user_id": self.user_id, "workspace_id": self.workspace_id, "active_flows": {k: v.to_dict() for k, v in self.active_flows.items()}, "active_agents": {k: v.to_dict() for k, v in self.active_agents.items()}, "interaction_history": self.interaction_history, "settings": self.settings, "context": self.context, "total_interactions": self.total_interactions, "total_flows_completed": self.total_flows_completed, "total_agents_completed": self.total_agents_completed, "created_at": to_iso(self.created_at), "last_activity_at": to_iso(self.last_activity_at), } @classmethod def from_dict(cls, data: dict[str, Any]) -> SessionState: """Создать из словаря.""" return cls( session_id=data["session_id"], user_id=data.get("user_id"), workspace_id=data.get("workspace_id", "default"), active_flows={ k: FlowState.from_dict(v) for k, v in data.get("active_flows", {}).items() }, active_agents={ k: AgentState.from_dict(v) for k, v in data.get("active_agents", {}).items() }, interaction_history=data.get("interaction_history", []), settings=data.get("settings", {}), context=data.get("context", {}), total_interactions=data.get("total_interactions", 0), total_flows_completed=data.get("total_flows_completed", 0), total_agents_completed=data.get("total_agents_completed", 0), created_at=parse_datetime(data.get("created_at")) or utc_now(), last_activity_at=parse_datetime(data.get("last_activity_at")) or utc_now(), ) # ============================================================================ # State Event # ============================================================================ @dataclass class StateEvent: """ Событие изменения состояния. Используется для отслеживания изменений и аудита. """ id: str = field(default_factory=lambda: str(uuid4())) event_type: str = "" # "state_changed", "agent_started", "flow_completed", etc. timestamp: datetime = field(default_factory=utc_now) # Контекст session_id: str | None = None flow_id: str | None = None agent_id: str | None = None # Данные old_value: Any = None new_value: Any = None metadata: dict[str, Any] = field(default_factory=dict) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "event_type": self.event_type, "timestamp": to_iso(self.timestamp), "session_id": self.session_id, "flow_id": self.flow_id, "agent_id": self.agent_id, "old_value": self.old_value, "new_value": self.new_value, "metadata": self.metadata, } # ============================================================================ # State Manager Protocol # ============================================================================ class StateManager(Protocol): """ Протокол для менеджера состояния. Определяет интерфейс для управления состоянием. """ async def save_session(self, session: SessionState) -> None: """Сохранить состояние сессии.""" ... async def load_session(self, session_id: str) -> SessionState | None: """Загрузить состояние сессии.""" ... async def delete_session(self, session_id: str) -> bool: """Удалить состояние сессии.""" ... async def list_sessions( self, user_id: str | None = None, limit: int = 100 ) -> list[SessionState]: """Получить список сессий.""" ... async def emit_event(self, event: StateEvent) -> None: """Отправить событие изменения состояния.""" ... async def subscribe(self, callback: Callable[[StateEvent], None]) -> None: """Подписаться на события.""" ... # ============================================================================ # Helper Functions # ============================================================================ def create_agent_state( agent_id: str, agent_name: str, input_data: dict[str, Any] | None = None, **kwargs: Any, ) -> AgentState: """Создать состояние агента.""" return AgentState( agent_id=agent_id, agent_name=agent_name, input_data=input_data or {}, **kwargs, ) def create_flow_state( flow_id: str, flow_name: str, **kwargs: Any, ) -> FlowState: """Создать состояние flow.""" return FlowState( flow_id=flow_id, flow_name=flow_name, **kwargs, ) def create_session_state( session_id: str, user_id: str | None = None, workspace_id: str = "default", **kwargs: Any, ) -> SessionState: """Создать состояние сессии.""" return SessionState( session_id=session_id, user_id=user_id, workspace_id=workspace_id, **kwargs, ) def create_execution_state( step_name: str, input_data: dict[str, Any] | None = None, **kwargs: Any, ) -> ExecutionState: """Создать состояние выполнения.""" return ExecutionState( step_name=step_name, input_data=input_data or {}, **kwargs, ) # ============================================================================ # Exports # ============================================================================ __all__ = [ "AgentState", "Checkpoint", "EdgeStatus", "ExecutionState", "FlowState", "NodeStatus", "RetryStrategy", "SessionState", "StateEvent", "StateManager", "TaskStatus", "create_agent_state", "create_execution_state", "create_flow_state", "create_session_state", ]