update
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user