From b54f5e23f3447165ac164483a4900bcf00cc503b Mon Sep 17 00:00:00 2001 From: team3 Date: Mon, 22 Jun 2026 06:07:22 +0200 Subject: [PATCH] update --- backend/agents.py | 21 +++++++++++++++++++++ backend/bausteine.py | 6 ++++-- backend/config.py | 2 +- backend/pipeline.py | 6 ++++-- 4 files changed, 30 insertions(+), 5 deletions(-) diff --git a/backend/agents.py b/backend/agents.py index 16281d7..139fce0 100644 --- a/backend/agents.py +++ b/backend/agents.py @@ -20,6 +20,23 @@ log = logging.getLogger("creator.agents") _active_processes: dict[str, asyncio.subprocess.Process] = {} +# Abgebrochene Scopes (Schlüssel-Präfixe, symmetrisch zu kill_process). Ein Agent, dessen +# Key mit einem dieser Präfixe beginnt, bricht VOR dem Spawn ab — so werden auch in der +# Semaphore-Schlange WARTENDE Agenten beim Abbruch sofort gestoppt, statt noch zu starten. +_cancelled_prefixes: set[str] = set() + + +def cancel_scope(prefix: str) -> None: + _cancelled_prefixes.add(prefix) + + +def clear_scope(prefix: str) -> None: + _cancelled_prefixes.discard(prefix) + + +def _scope_cancelled(agent_key: str) -> bool: + return any(agent_key.startswith(p) for p in _cancelled_prefixes) + # Deckelt die realen CLI-Prozesse — unabhängig von der Pipeline-Semaphore in # generator.py. Acquire passiert VOR dem Spawn, damit Wartezeit in der Queue # nicht gegen den Agent-Timeout zählt. @@ -87,12 +104,16 @@ async def run_agent( capabilities: str = "none", lane: str = "batch", ) -> tuple[int, str, str]: + if _scope_cancelled(agent_key): # vor dem Anstehen: gar nicht erst in die Schlange + return 1, "", "abgebrochen" if provider not in PROVIDERS: return 1, "", f"Unbekannter Provider: {provider}" if shutil.which(PROVIDERS[provider]["cli"]) is None: return 1, "", f"CLI '{PROVIDERS[provider]['cli']}' nicht installiert (Provider: {provider})" sem = _interactive_sem if lane == "interactive" else _batch_sem async with sem: + if _scope_cancelled(agent_key): # nach dem Acquire: in der Schlange abgebrochen → kein Spawn + return 1, "", "abgebrochen" if PROVIDERS[provider]["cli"] == "opencode": return await _run_opencode(agent_key, prompt, timeout, provider, role, capabilities) return await _run_claude_cli(agent_key, prompt, timeout, role, capabilities) diff --git a/backend/bausteine.py b/backend/bausteine.py index 7e4ad14..7971556 100644 --- a/backend/bausteine.py +++ b/backend/bausteine.py @@ -18,7 +18,7 @@ import shutil import subprocess from pathlib import Path -from agents import kill_process +from agents import kill_process, cancel_scope, clear_scope from config import KONSENS_GRACE, KONSENS_MAX_RUNDEN, DEFAULT_PROVIDER from fsutil import atomic_write_text, atomic_write_json from jsonio import read_json_file as _json_datei @@ -243,7 +243,8 @@ def cancel_bausteine(topic: str) -> bool: if topic not in _bausteine_progress: return False _bausteine_cancelled.add(topic) - kill_process(f"bausteine-{topic}-") + cancel_scope(f"bausteine-{topic}-") # wartende Agenten bailen vorm Spawn + kill_process(f"bausteine-{topic}-") # laufende Subprozesse killen return True @@ -1473,3 +1474,4 @@ async def generate_bausteine(topic: str, instructions: str = "", provider: str = _bausteine_progress.pop(topic, None) _bausteine_step.pop(topic, None) _bausteine_cancelled.discard(topic) + clear_scope(f"bausteine-{topic}-") # Scope leeren → Neustart blockiert nicht diff --git a/backend/config.py b/backend/config.py index dff871c..d5849f2 100644 --- a/backend/config.py +++ b/backend/config.py @@ -24,7 +24,7 @@ LESBARKEIT_HART_ANTEIL = 0.30 # … ODER wenn dieser Anteil der Sätze hart ist # Deckel für gleichzeitige CLI-Agenten-Prozesse (über alle Generierungen hinweg). # Eigene Spur für interaktive Aufrufe (Chat, Elemente), damit sie nicht hinter # laufenden Writern in der Warteschlange hängen. -MAX_CONCURRENT_AGENTS = 12 +MAX_CONCURRENT_AGENTS = 20 MAX_CONCURRENT_INTERACTIVE = 4 # Grace-Fenster der Konsens-Races (Bausteine, Guide, OnePager): Nach dem ersten diff --git a/backend/pipeline.py b/backend/pipeline.py index 36e48d3..777854d 100644 --- a/backend/pipeline.py +++ b/backend/pipeline.py @@ -12,7 +12,7 @@ from datetime import datetime, timezone from pathlib import Path from typing import Callable -from agents import run_agent, kill_process +from agents import run_agent, kill_process, cancel_scope, clear_scope from config import MAX_CONCURRENT_GENERATIONS, TEMPLATES_DIR, TIMEOUTS from database import update_guide from jsonio import read_json_file as _json_datei @@ -26,7 +26,8 @@ _cancelled: set[str] = set() async def cancel_guide(guide_id: str) -> bool: _cancelled.add(guide_id) - kill_process(guide_id) + cancel_scope(f"{guide_id}-") # wartende Agenten bailen vorm Spawn + kill_process(guide_id) # laufende Subprozesse killen now = datetime.now(timezone.utc).isoformat() await update_guide(guide_id, status="error", progress=None, error_msg="Abgebrochen — Fortschritt bleibt erhalten", updated_at=now) return True @@ -38,6 +39,7 @@ def is_guide_cancelled(guide_id: str) -> bool: def clear_guide_cancelled(guide_id: str) -> None: _cancelled.discard(guide_id) + clear_scope(f"{guide_id}-") # Scope leeren → Neustart blockiert nicht async def _set_progress(guide_id: str, progress: str) -> None: