Elevate MiroFish/CrowdSight from single-container dev to a SaaS foundation: - Local memory backend (Zep-compatible): memory services/models, local graph builder + updater, AgentActivity seam, import-boundary isolation; Zep stays default, local is opt-in behind MEMORY_BACKEND. Semantic parity not yet proven. - Durable product persistence: projects/simulations/reports schema (migration 0007) + tenant/owner-scoped ProductRepository + dual-write + scoped_project read-first + ArtifactStore abstraction; durable JobQueue + worker.py. - SaaS hardening: durable RateLimiter (wired to login), UsageService (LLM accounting), redacted AuditService, idempotency, CORS allowlist, safe API errors, single-use PasswordResetService + endpoints (covers invite-pending). - Exactly 3 roles (super_admin/admin/user) with tenant authz policy. - Admin UI: GET/POST/PATCH /api/admin/users + GET/PUT /api/admin/settings (super-admin only, encrypted/masked); AdminView.vue + SettingsView.vue with admin/super-admin route guards, th/en i18n. - Production deploy topology: multi-stage Dockerfile (frontend build + gunicorn wsgi + nginx SPA-proxy + supervisord worker), backend/wsgi.py, gunicorn dep. Backend 197 passed; frontend 10 tests + build green. ruff unavailable (gap). No commit of credentials; secrets handled via env/.env.example. Deferred: Zep semantic A/B parity, object storage cutover, mobile QA, EasyPanel container build of deploy topology.
178 lines
5.7 KiB
Python
178 lines
5.7 KiB
Python
"""Tenant-scoped local persistence for simulation activities.
|
|
|
|
This module intentionally does not import or construct a Zep client. It keeps
|
|
runtime activity updates behind the same small manager interface used by the
|
|
legacy runner, while storing each activity as a durable local memory episode.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import threading
|
|
from datetime import datetime
|
|
from typing import Any, Dict, Optional, cast
|
|
|
|
from sqlalchemy.orm import Session
|
|
|
|
from .memory_repository import SqlAlchemyMemoryRepository
|
|
from ..utils.logger import get_logger
|
|
|
|
logger = get_logger("crowdsight.local_graph_memory_updater")
|
|
|
|
|
|
class LocalGraphMemoryUpdater:
|
|
"""Persist simulation activities in a tenant-scoped local graph."""
|
|
|
|
def __init__(
|
|
self,
|
|
simulation_id: str,
|
|
graph_id: str,
|
|
organization_id: str,
|
|
session_factory,
|
|
):
|
|
for name, value in (
|
|
("simulation_id", simulation_id),
|
|
("graph_id", graph_id),
|
|
("organization_id", organization_id),
|
|
):
|
|
if not isinstance(value, str) or not value.strip():
|
|
raise ValueError(f"local_memory_{name}_required")
|
|
if not callable(session_factory):
|
|
raise ValueError("memory_session_factory_required")
|
|
|
|
self.simulation_id = simulation_id
|
|
self.graph_id = graph_id
|
|
self.organization_id = organization_id
|
|
self.session_factory = session_factory
|
|
self._running = False
|
|
self._sequence = 0
|
|
self._lock = threading.Lock()
|
|
self._total_activities = 0
|
|
self._total_sent = 0
|
|
self._total_items_sent = 0
|
|
self._failed_count = 0
|
|
self._skipped_count = 0
|
|
|
|
def start(self) -> None:
|
|
self._running = True
|
|
|
|
def stop(self) -> None:
|
|
self._running = False
|
|
|
|
def add_activity(self, activity: Any) -> None:
|
|
if not callable(getattr(activity, "to_episode_text", None)):
|
|
raise ValueError("invalid_agent_activity")
|
|
if activity.action_type == "DO_NOTHING":
|
|
self._skipped_count += 1
|
|
return
|
|
|
|
with self._lock:
|
|
self._sequence += 1
|
|
sequence = self._sequence
|
|
self._total_activities += 1
|
|
|
|
try:
|
|
episode_text = activity.to_episode_text()
|
|
source_ref = (
|
|
f"{self.simulation_id}:{activity.platform}:{activity.round_num}:"
|
|
f"{activity.agent_id}:{sequence}"
|
|
)
|
|
session: Session = cast(Session, self.session_factory())
|
|
try:
|
|
repository = SqlAlchemyMemoryRepository(
|
|
session,
|
|
organization_id=self.organization_id,
|
|
graph_id=self.graph_id,
|
|
)
|
|
repository.add_episode(
|
|
source_type="simulation_activity",
|
|
source_ref=source_ref,
|
|
normalized_text=episode_text,
|
|
summary=episode_text,
|
|
extractor_version="runtime-v1",
|
|
)
|
|
session.commit()
|
|
finally:
|
|
session.close()
|
|
self._total_sent += 1
|
|
self._total_items_sent += 1
|
|
except Exception:
|
|
self._failed_count += 1
|
|
raise
|
|
|
|
def add_activity_from_dict(self, data: Dict[str, Any], platform: str) -> None:
|
|
if "event_type" in data:
|
|
return
|
|
from .memory_activity import AgentActivity
|
|
|
|
self.add_activity(
|
|
AgentActivity(
|
|
platform=platform,
|
|
agent_id=data.get("agent_id", 0),
|
|
agent_name=data.get("agent_name", ""),
|
|
action_type=data.get("action_type", ""),
|
|
action_args=data.get("action_args", {}),
|
|
round_num=data.get("round", 0),
|
|
timestamp=data.get("timestamp", datetime.now().isoformat()),
|
|
)
|
|
)
|
|
|
|
def get_stats(self) -> Dict[str, Any]:
|
|
return {
|
|
"total_activities": self._total_activities,
|
|
"batches_sent": self._total_sent,
|
|
"items_sent": self._total_items_sent,
|
|
"failed_count": self._failed_count,
|
|
"skipped_count": self._skipped_count,
|
|
"queue_size": 0,
|
|
"buffer_sizes": {},
|
|
"running": self._running,
|
|
}
|
|
|
|
|
|
class LocalGraphMemoryManager:
|
|
"""Manage one local updater per simulation."""
|
|
|
|
_updaters: Dict[str, LocalGraphMemoryUpdater] = {}
|
|
_lock = threading.Lock()
|
|
|
|
@classmethod
|
|
def create_updater(
|
|
cls,
|
|
simulation_id: str,
|
|
graph_id: str,
|
|
*,
|
|
organization_id: str,
|
|
session_factory,
|
|
) -> LocalGraphMemoryUpdater:
|
|
with cls._lock:
|
|
existing = cls._updaters.get(simulation_id)
|
|
if existing is not None:
|
|
existing.stop()
|
|
updater = LocalGraphMemoryUpdater(
|
|
simulation_id=simulation_id,
|
|
graph_id=graph_id,
|
|
organization_id=organization_id,
|
|
session_factory=session_factory,
|
|
)
|
|
updater.start()
|
|
cls._updaters[simulation_id] = updater
|
|
return updater
|
|
|
|
@classmethod
|
|
def get_updater(cls, simulation_id: str) -> Optional[LocalGraphMemoryUpdater]:
|
|
return cls._updaters.get(simulation_id)
|
|
|
|
@classmethod
|
|
def stop_updater(cls, simulation_id: str) -> None:
|
|
with cls._lock:
|
|
updater = cls._updaters.pop(simulation_id, None)
|
|
if updater is not None:
|
|
updater.stop()
|
|
|
|
@classmethod
|
|
def stop_all(cls) -> None:
|
|
with cls._lock:
|
|
simulation_ids = list(cls._updaters)
|
|
for simulation_id in simulation_ids:
|
|
cls.stop_updater(simulation_id)
|