Files
microfish/backend/app/services/local_graph_memory_updater.py
Kunthawat Greethong 8b84378fe1 feat: SaaS foundation for CrowdSight
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.
2026-08-31 13:05:21 +07:00

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)