2 Commits

Autor SHA1 Nachricht Datum
Claude Code
3a3076c4cf 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>
2026-07-23 20:12:03 +00:00
e0107a1bb1 Merge pull request 'Hotfix: E-Mail-Versand mit Retry (IONOS-Farm-Stoerung)' (#50) from fix/email-retry-live into main 2026-07-10 12:20:21 +02:00
2 geänderte Dateien mit 123 neuen und 65 gelöschten Zeilen

Datei anzeigen

@@ -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

Datei anzeigen

@@ -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,