Orchestrator: eigener Worker pro Organisation (parallele Lanes)
Ersetzt die globale Warteschlange mit einem einzelnen Worker durch Per-Tenant-Lanes: jede Organisation bekommt eine eigene Queue und einen eigenen Worker-Task, die unabhaengig und parallel laufen. Lanes werden lazy angelegt und bei Leerlauf (60s) wieder beendet. - Skalarer Zustand (_current_task, _cancel_event) -> Maps pro incident_id (_current_tasks, _cancel_events); _queued_ids/_cancel_requested bleiben global (IDs sind lane-uebergreifend eindeutig). - enqueue_refresh/cancel_refresh: Signatur unveraendert, routen intern per tenant_id in die passende Lane. Oeffentliche Lagen (tenant_id NULL) in fester PUBLIC_LANE. - Optionales Ventil ORCHESTRATOR_MAX_PARALLEL (Default 0 = unbegrenzt) begrenzt bei Bedarf gleichzeitige Recherchen ueber alle Lanes. - /refreshing-Endpoint auf _current_tasks-Map umgestellt, Response-Format unveraendert (pro Tenant laeuft weiterhin hoechstens eine Lage). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Dieser Commit ist enthalten in:
@@ -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
|
||||
|
||||
In neuem Issue referenzieren
Einen Benutzer sperren