diff --git a/src/agents/orchestrator.py b/src/agents/orchestrator.py index 86bf7e8..6cfd306 100644 --- a/src/agents/orchestrator.py +++ b/src/agents/orchestrator.py @@ -2,6 +2,7 @@ import asyncio import json import logging +import os import re from datetime import datetime from config import TIMEZONE @@ -401,48 +402,83 @@ async def _send_email_notifications_for_incident( class AgentOrchestrator: - """Verwaltet die Claude-Agenten-Queue und koordiniert Recherche-Zyklen.""" + """Koordiniert die Recherche-Zyklen: pro Organisation eine unabhaengige, + parallel laufende Worker-Lane.""" + + # Lane-Schluessel fuer oeffentliche/systemweite Lagen (tenant_id IS NULL). + PUBLIC_LANE = 0 + # Sekunden ohne Auftrag, nach denen eine Lane sich beendet und aufraeumt. + IDLE_TIMEOUT = 60 def __init__(self): - self._queue: asyncio.Queue = asyncio.Queue() self._running = False - self._current_task: Optional[int] = None - # Session-Start des aktuellen Tasks (UTC ISO mit 'Z'). Ueberspannt Multi-Pass - # und Retries innerhalb derselben Queue-Abarbeitung — verhindert, dass der - # Frontend-Timer beim Seiten-Reload auf den Pass/Retry-Start zurueckspringt. - self._current_task_started_at: Optional[str] = None self._ws_manager = None - self._queued_ids: set[int] = set() + # Pro Organisation (tenant_id) eine eigene Queue + ein eigener Worker-Task. + # So arbeiten Organisationen unabhaengig und parallel, statt sich eine + # globale Warteschlange zu teilen. Lanes werden bei Bedarf angelegt und bei + # Leerlauf wieder beendet. Der Lock schuetzt das Anlegen/Beenden gegen + # gleichzeitiges Einreihen. + self._lanes: dict[int, asyncio.Queue] = {} + self._lane_workers: dict[int, asyncio.Task] = {} + self._lanes_lock = asyncio.Lock() + # Globale Zustaende — incident-IDs sind lane-uebergreifend eindeutig. + self._queued_ids: set[int] = set() # eingereiht (in irgendeiner Lane) + # incident_id -> Session-Start (UTC ISO mit 'Z'). Ersetzt den frueheren + # Einzelwert; haelt den Frontend-Timer ueber Multi-Pass/Retry stabil. + self._current_tasks: dict[int, str] = {} self._cancel_requested: set[int] = set() - self._cancel_event: asyncio.Event | None = None + # incident_id -> Cancel-Event des laufenden Tasks. cancel_refresh laeuft in + # einem anderen async-Kontext als der Worker und kann die ContextVar nicht + # nutzen, daher diese direkte Zuordnung. + self._cancel_events: dict[int, asyncio.Event] = {} + # Optionales Sicherheitsventil gegen Ueberlastung des gemeinsamen Claude- + # Kontos: begrenzt die Zahl GLEICHZEITIG laufender Recherchen ueber alle + # Lanes. Default 0 = unbegrenzt (jede Organisation voellig unabhaengig). + _max = int(os.getenv("ORCHESTRATOR_MAX_PARALLEL", "0") or "0") + self._global_sem: Optional[asyncio.Semaphore] = asyncio.Semaphore(_max) if _max > 0 else None def set_ws_manager(self, ws_manager): """WebSocket-Manager setzen für Echtzeit-Updates.""" self._ws_manager = ws_manager async def start(self): - """Queue-Worker starten.""" + """Orchestrator aktivieren. Worker-Lanes werden pro Organisation bei der + ersten Anfrage angelegt (lazy).""" self._running = True - asyncio.create_task(self._worker()) logger.info("Agenten-Orchestrator gestartet") async def stop(self): - """Queue-Worker stoppen.""" + """Orchestrator stoppen und alle Lane-Worker beenden.""" self._running = False + async with self._lanes_lock: + workers = list(self._lane_workers.values()) + self._lane_workers.clear() + self._lanes.clear() + for task in workers: + task.cancel() logger.info("Agenten-Orchestrator gestoppt") async def enqueue_refresh(self, incident_id: int, trigger_type: str = "manual", user_id: int = None) -> bool: - """Refresh-Auftrag in die Queue stellen. Gibt False zurueck wenn bereits in Queue/aktiv.""" - if incident_id in self._queued_ids or self._current_task == incident_id: + """Refresh-Auftrag in die Lane der Organisation stellen. Gibt False zurueck + wenn die Lage dort bereits wartet oder gerade laeuft.""" + if incident_id in self._queued_ids or incident_id in self._current_tasks: logger.info(f"Refresh fuer Lage {incident_id} uebersprungen: bereits aktiv/in Queue") return False visibility, created_by, tenant_id = await self._get_incident_visibility(incident_id) + lane_key = tenant_id if tenant_id else self.PUBLIC_LANE - self._queued_ids.add(incident_id) - await self._queue.put((incident_id, trigger_type, user_id)) - queue_size = self._queue.qsize() - logger.info(f"Refresh fuer Lage {incident_id} eingereiht (Queue: {queue_size}, Trigger: {trigger_type})") + async with self._lanes_lock: + queue = self._lanes.get(lane_key) + if queue is None: + queue = asyncio.Queue() + self._lanes[lane_key] = queue + self._lane_workers[lane_key] = asyncio.create_task(self._worker(lane_key, queue)) + logger.info(f"Neue Worker-Lane fuer Organisation {lane_key} gestartet") + self._queued_ids.add(incident_id) + queue.put_nowait((incident_id, trigger_type, user_id)) + queue_size = queue.qsize() + logger.info(f"Refresh fuer Lage {incident_id} eingereiht (Lane {lane_key}, Queue: {queue_size}, Trigger: {trigger_type})") if self._ws_manager: await self._ws_manager.broadcast_for_incident({ @@ -455,11 +491,12 @@ class AgentOrchestrator: async def cancel_refresh(self, incident_id: int) -> bool: """Fordert Abbruch eines laufenden oder wartenden Refreshes an.""" - # Check if it's the currently running task - if self._current_task == incident_id: + # Laeuft die Lage gerade? + if incident_id in self._current_tasks: self._cancel_requested.add(incident_id) - if self._cancel_event: - self._cancel_event.set() + ev = self._cancel_events.get(incident_id) + if ev: + ev.set() logger.info(f"Cancel angefordert fuer laufende Lage {incident_id}") if self._ws_manager: try: @@ -473,25 +510,33 @@ class AgentOrchestrator: }, vis, cb, tid) return True - # Check if it's in the queue (not yet started) + # Wartet die Lage noch in ihrer Lane? if incident_id in self._queued_ids: self._queued_ids.discard(incident_id) - # Remove from asyncio queue (rebuild without this ID) + # Betroffene Lane bestimmen und die ID aus deren Queue entfernen. + try: + _vis, _cb, _tid = await self._get_incident_visibility(incident_id) + except Exception: + _tid = None + lane_key = _tid if _tid else self.PUBLIC_LANE removed = False - new_items = [] - while not self._queue.empty(): - try: - item = self._queue.get_nowait() - iid = item[0] if isinstance(item, tuple) else item - if iid == incident_id: - removed = True - self._queue.task_done() - else: - new_items.append(item) - except Exception: - break - for item in new_items: - self._queue.put_nowait(item) + async with self._lanes_lock: + queue = self._lanes.get(lane_key) + if queue is not None: + new_items = [] + while not queue.empty(): + try: + item = queue.get_nowait() + except Exception: + break + iid = item[0] if isinstance(item, tuple) else item + if iid == incident_id: + removed = True + queue.task_done() + else: + new_items.append(item) + for item in new_items: + queue.put_nowait(item) logger.info(f"Lage {incident_id} aus Warteschlange entfernt (removed={removed})") @@ -519,12 +564,22 @@ class AgentOrchestrator: self._cancel_requested.discard(incident_id) raise asyncio.CancelledError("Vom Nutzer abgebrochen") - async def _worker(self): - """Verarbeitet Refresh-Aufträge sequentiell.""" + async def _worker(self, lane_key: int, queue: asyncio.Queue): + """Verarbeitet die Auftraege EINER Organisation sequentiell. Verschiedene + Organisationen laufen in eigenen Lanes parallel. Bei Leerlauf beendet sich + die Lane selbst und wird beim naechsten Auftrag neu angelegt.""" while self._running: try: - item = await asyncio.wait_for(self._queue.get(), timeout=5.0) + item = await asyncio.wait_for(queue.get(), timeout=self.IDLE_TIMEOUT) except asyncio.TimeoutError: + # Leerlauf: Lane beenden. Unter Lock gegen gleichzeitiges Einreihen, + # damit kein Auftrag in einer verwaisten Queue liegen bleibt. + async with self._lanes_lock: + if queue.empty() and self._lanes.get(lane_key) is queue: + self._lanes.pop(lane_key, None) + self._lane_workers.pop(lane_key, None) + logger.info(f"Worker-Lane fuer Organisation {lane_key} bei Leerlauf beendet") + return continue if len(item) == 3: @@ -533,12 +588,12 @@ class AgentOrchestrator: incident_id, trigger_type = item user_id = None self._queued_ids.discard(incident_id) - self._current_task = incident_id # Session-Start EINMAL setzen — bleibt ueber Multi-Pass/Retry hinweg stabil - self._current_task_started_at = datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%SZ') - self._cancel_event = asyncio.Event() - _cancel_event_var.set(self._cancel_event) - logger.info(f"Starte Refresh für Lage {incident_id} (Trigger: {trigger_type})") + self._current_tasks[incident_id] = datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%SZ') + cancel_event = asyncio.Event() + self._cancel_events[incident_id] = cancel_event + _cancel_event_var.set(cancel_event) + logger.info(f"Starte Refresh für Lage {incident_id} (Lane {lane_key}, Trigger: {trigger_type})") RETRY_DELAYS = [0, 120, 300] # Sekunden: sofort, 2min, 5min TRANSIENT_ERRORS = (asyncio.TimeoutError, TimeoutError, ConnectionError, OSError) @@ -548,6 +603,10 @@ class AgentOrchestrator: def _is_transient_cli(err: Exception) -> bool: return isinstance(err, ClaudeCliError) and err.error_type in ("rate_limit", "timeout") + # Optionales globales Ventil (Default aus): begrenzt gleichzeitige Recherchen. + sem = self._global_sem + if sem is not None: + await sem.acquire() try: # Research-Lagen: Automatisch 3 Durchläufe nur beim ersten Refresh incident_type, has_summary = await self._get_incident_info(incident_id) @@ -626,11 +685,13 @@ class AgentOrchestrator: "data": {"error": str(last_error)}, }, _vis, _cb, _tid) finally: - self._current_task = None - self._current_task_started_at = None - self._cancel_event = None + if sem is not None: + sem.release() + self._current_tasks.pop(incident_id, None) + self._cancel_events.pop(incident_id, None) + self._cancel_requested.discard(incident_id) _cancel_event_var.set(None) - self._queue.task_done() + queue.task_done() async def _mark_refresh_cancelled(self, incident_id: int): """Markiert den laufenden Refresh-Log-Eintrag als cancelled und schliesst diff --git a/src/routers/incidents.py b/src/routers/incidents.py index 4760a4f..f9c6144 100644 --- a/src/routers/incidents.py +++ b/src/routers/incidents.py @@ -165,28 +165,25 @@ async def get_refreshing_incidents( ) rows = await cursor.fetchall() - # Also include queued incidents from orchestrator + # Queued- und laufende Lagen aus dem Orchestrator ergaenzen. from agents.orchestrator import orchestrator - queued_ids = list(orchestrator._queued_ids) if hasattr(orchestrator, '_queued_ids') else [] - current_task = orchestrator._current_task if hasattr(orchestrator, '_current_task') else None - # Session-Start des aktuell laufenden Tasks — stabil ueber Multi-Pass/Retry hinweg. - # Verhindert, dass der Frontend-Timer beim Reload auf den letzten Log-Eintrag - # (pass 2/3 oder retry n) zurueckspringt. - current_started_at = ( - orchestrator._current_task_started_at - if hasattr(orchestrator, '_current_task_started_at') else None - ) + queued_ids = list(getattr(orchestrator, '_queued_ids', set())) + # incident_id -> stabiler Session-Start (ueber Multi-Pass/Retry hinweg). Ersetzt + # den frueheren Einzelwert, da jetzt mehrere Organisationen parallel laufen. + current_tasks = getattr(orchestrator, '_current_tasks', {}) or {} details = {} for row in rows: iid = row["incident_id"] - started_at = ( - current_started_at - if (iid == current_task and current_started_at) - else row["started_at"] - ) + session_start = current_tasks.get(iid) + started_at = session_start if session_start else row["started_at"] details[str(iid)] = {"started_at": started_at} + # Pro Organisation laeuft hoechstens eine Lage gleichzeitig; der Endpoint ist + # ohnehin tenant-gefiltert, daher genuegt der erste laufende Treffer. + running_here = [row["incident_id"] for row in rows if row["incident_id"] in current_tasks] + current_task = running_here[0] if running_here else None + return { "refreshing": [row["incident_id"] for row in rows], "queued": queued_ids,