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 asyncio
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
|
import os
|
||||||
import re
|
import re
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from config import TIMEZONE
|
from config import TIMEZONE
|
||||||
@@ -401,48 +402,83 @@ async def _send_email_notifications_for_incident(
|
|||||||
|
|
||||||
|
|
||||||
class AgentOrchestrator:
|
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):
|
def __init__(self):
|
||||||
self._queue: asyncio.Queue = asyncio.Queue()
|
|
||||||
self._running = False
|
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._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_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):
|
def set_ws_manager(self, ws_manager):
|
||||||
"""WebSocket-Manager setzen für Echtzeit-Updates."""
|
"""WebSocket-Manager setzen für Echtzeit-Updates."""
|
||||||
self._ws_manager = ws_manager
|
self._ws_manager = ws_manager
|
||||||
|
|
||||||
async def start(self):
|
async def start(self):
|
||||||
"""Queue-Worker starten."""
|
"""Orchestrator aktivieren. Worker-Lanes werden pro Organisation bei der
|
||||||
|
ersten Anfrage angelegt (lazy)."""
|
||||||
self._running = True
|
self._running = True
|
||||||
asyncio.create_task(self._worker())
|
|
||||||
logger.info("Agenten-Orchestrator gestartet")
|
logger.info("Agenten-Orchestrator gestartet")
|
||||||
|
|
||||||
async def stop(self):
|
async def stop(self):
|
||||||
"""Queue-Worker stoppen."""
|
"""Orchestrator stoppen und alle Lane-Worker beenden."""
|
||||||
self._running = False
|
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")
|
logger.info("Agenten-Orchestrator gestoppt")
|
||||||
|
|
||||||
async def enqueue_refresh(self, incident_id: int, trigger_type: str = "manual", user_id: int = None) -> bool:
|
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."""
|
"""Refresh-Auftrag in die Lane der Organisation stellen. Gibt False zurueck
|
||||||
if incident_id in self._queued_ids or self._current_task == incident_id:
|
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")
|
logger.info(f"Refresh fuer Lage {incident_id} uebersprungen: bereits aktiv/in Queue")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
visibility, created_by, tenant_id = await self._get_incident_visibility(incident_id)
|
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)
|
async with self._lanes_lock:
|
||||||
await self._queue.put((incident_id, trigger_type, user_id))
|
queue = self._lanes.get(lane_key)
|
||||||
queue_size = self._queue.qsize()
|
if queue is None:
|
||||||
logger.info(f"Refresh fuer Lage {incident_id} eingereiht (Queue: {queue_size}, Trigger: {trigger_type})")
|
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:
|
if self._ws_manager:
|
||||||
await self._ws_manager.broadcast_for_incident({
|
await self._ws_manager.broadcast_for_incident({
|
||||||
@@ -455,11 +491,12 @@ class AgentOrchestrator:
|
|||||||
|
|
||||||
async def cancel_refresh(self, incident_id: int) -> bool:
|
async def cancel_refresh(self, incident_id: int) -> bool:
|
||||||
"""Fordert Abbruch eines laufenden oder wartenden Refreshes an."""
|
"""Fordert Abbruch eines laufenden oder wartenden Refreshes an."""
|
||||||
# Check if it's the currently running task
|
# Laeuft die Lage gerade?
|
||||||
if self._current_task == incident_id:
|
if incident_id in self._current_tasks:
|
||||||
self._cancel_requested.add(incident_id)
|
self._cancel_requested.add(incident_id)
|
||||||
if self._cancel_event:
|
ev = self._cancel_events.get(incident_id)
|
||||||
self._cancel_event.set()
|
if ev:
|
||||||
|
ev.set()
|
||||||
logger.info(f"Cancel angefordert fuer laufende Lage {incident_id}")
|
logger.info(f"Cancel angefordert fuer laufende Lage {incident_id}")
|
||||||
if self._ws_manager:
|
if self._ws_manager:
|
||||||
try:
|
try:
|
||||||
@@ -473,25 +510,33 @@ class AgentOrchestrator:
|
|||||||
}, vis, cb, tid)
|
}, vis, cb, tid)
|
||||||
return True
|
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:
|
if incident_id in self._queued_ids:
|
||||||
self._queued_ids.discard(incident_id)
|
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
|
removed = False
|
||||||
new_items = []
|
async with self._lanes_lock:
|
||||||
while not self._queue.empty():
|
queue = self._lanes.get(lane_key)
|
||||||
try:
|
if queue is not None:
|
||||||
item = self._queue.get_nowait()
|
new_items = []
|
||||||
iid = item[0] if isinstance(item, tuple) else item
|
while not queue.empty():
|
||||||
if iid == incident_id:
|
try:
|
||||||
removed = True
|
item = queue.get_nowait()
|
||||||
self._queue.task_done()
|
except Exception:
|
||||||
else:
|
break
|
||||||
new_items.append(item)
|
iid = item[0] if isinstance(item, tuple) else item
|
||||||
except Exception:
|
if iid == incident_id:
|
||||||
break
|
removed = True
|
||||||
for item in new_items:
|
queue.task_done()
|
||||||
self._queue.put_nowait(item)
|
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})")
|
logger.info(f"Lage {incident_id} aus Warteschlange entfernt (removed={removed})")
|
||||||
|
|
||||||
@@ -519,12 +564,22 @@ class AgentOrchestrator:
|
|||||||
self._cancel_requested.discard(incident_id)
|
self._cancel_requested.discard(incident_id)
|
||||||
raise asyncio.CancelledError("Vom Nutzer abgebrochen")
|
raise asyncio.CancelledError("Vom Nutzer abgebrochen")
|
||||||
|
|
||||||
async def _worker(self):
|
async def _worker(self, lane_key: int, queue: asyncio.Queue):
|
||||||
"""Verarbeitet Refresh-Aufträge sequentiell."""
|
"""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:
|
while self._running:
|
||||||
try:
|
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:
|
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
|
continue
|
||||||
|
|
||||||
if len(item) == 3:
|
if len(item) == 3:
|
||||||
@@ -533,12 +588,12 @@ class AgentOrchestrator:
|
|||||||
incident_id, trigger_type = item
|
incident_id, trigger_type = item
|
||||||
user_id = None
|
user_id = None
|
||||||
self._queued_ids.discard(incident_id)
|
self._queued_ids.discard(incident_id)
|
||||||
self._current_task = incident_id
|
|
||||||
# Session-Start EINMAL setzen — bleibt ueber Multi-Pass/Retry hinweg stabil
|
# 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._current_tasks[incident_id] = datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%SZ')
|
||||||
self._cancel_event = asyncio.Event()
|
cancel_event = asyncio.Event()
|
||||||
_cancel_event_var.set(self._cancel_event)
|
self._cancel_events[incident_id] = cancel_event
|
||||||
logger.info(f"Starte Refresh für Lage {incident_id} (Trigger: {trigger_type})")
|
_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
|
RETRY_DELAYS = [0, 120, 300] # Sekunden: sofort, 2min, 5min
|
||||||
TRANSIENT_ERRORS = (asyncio.TimeoutError, TimeoutError, ConnectionError, OSError)
|
TRANSIENT_ERRORS = (asyncio.TimeoutError, TimeoutError, ConnectionError, OSError)
|
||||||
@@ -548,6 +603,10 @@ class AgentOrchestrator:
|
|||||||
def _is_transient_cli(err: Exception) -> bool:
|
def _is_transient_cli(err: Exception) -> bool:
|
||||||
return isinstance(err, ClaudeCliError) and err.error_type in ("rate_limit", "timeout")
|
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:
|
try:
|
||||||
# Research-Lagen: Automatisch 3 Durchläufe nur beim ersten Refresh
|
# Research-Lagen: Automatisch 3 Durchläufe nur beim ersten Refresh
|
||||||
incident_type, has_summary = await self._get_incident_info(incident_id)
|
incident_type, has_summary = await self._get_incident_info(incident_id)
|
||||||
@@ -626,11 +685,13 @@ class AgentOrchestrator:
|
|||||||
"data": {"error": str(last_error)},
|
"data": {"error": str(last_error)},
|
||||||
}, _vis, _cb, _tid)
|
}, _vis, _cb, _tid)
|
||||||
finally:
|
finally:
|
||||||
self._current_task = None
|
if sem is not None:
|
||||||
self._current_task_started_at = None
|
sem.release()
|
||||||
self._cancel_event = None
|
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)
|
_cancel_event_var.set(None)
|
||||||
self._queue.task_done()
|
queue.task_done()
|
||||||
|
|
||||||
async def _mark_refresh_cancelled(self, incident_id: int):
|
async def _mark_refresh_cancelled(self, incident_id: int):
|
||||||
"""Markiert den laufenden Refresh-Log-Eintrag als cancelled und schliesst
|
"""Markiert den laufenden Refresh-Log-Eintrag als cancelled und schliesst
|
||||||
|
|||||||
@@ -165,28 +165,25 @@ async def get_refreshing_incidents(
|
|||||||
)
|
)
|
||||||
rows = await cursor.fetchall()
|
rows = await cursor.fetchall()
|
||||||
|
|
||||||
# Also include queued incidents from orchestrator
|
# Queued- und laufende Lagen aus dem Orchestrator ergaenzen.
|
||||||
from agents.orchestrator import orchestrator
|
from agents.orchestrator import orchestrator
|
||||||
queued_ids = list(orchestrator._queued_ids) if hasattr(orchestrator, '_queued_ids') else []
|
queued_ids = list(getattr(orchestrator, '_queued_ids', set()))
|
||||||
current_task = orchestrator._current_task if hasattr(orchestrator, '_current_task') else None
|
# incident_id -> stabiler Session-Start (ueber Multi-Pass/Retry hinweg). Ersetzt
|
||||||
# Session-Start des aktuell laufenden Tasks — stabil ueber Multi-Pass/Retry hinweg.
|
# den frueheren Einzelwert, da jetzt mehrere Organisationen parallel laufen.
|
||||||
# Verhindert, dass der Frontend-Timer beim Reload auf den letzten Log-Eintrag
|
current_tasks = getattr(orchestrator, '_current_tasks', {}) or {}
|
||||||
# (pass 2/3 oder retry n) zurueckspringt.
|
|
||||||
current_started_at = (
|
|
||||||
orchestrator._current_task_started_at
|
|
||||||
if hasattr(orchestrator, '_current_task_started_at') else None
|
|
||||||
)
|
|
||||||
|
|
||||||
details = {}
|
details = {}
|
||||||
for row in rows:
|
for row in rows:
|
||||||
iid = row["incident_id"]
|
iid = row["incident_id"]
|
||||||
started_at = (
|
session_start = current_tasks.get(iid)
|
||||||
current_started_at
|
started_at = session_start if session_start else row["started_at"]
|
||||||
if (iid == current_task and current_started_at)
|
|
||||||
else row["started_at"]
|
|
||||||
)
|
|
||||||
details[str(iid)] = {"started_at": 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 {
|
return {
|
||||||
"refreshing": [row["incident_id"] for row in rows],
|
"refreshing": [row["incident_id"] for row in rows],
|
||||||
"queued": queued_ids,
|
"queued": queued_ids,
|
||||||
|
|||||||
In neuem Issue referenzieren
Einen Benutzer sperren