"""Provider-Schicht: ein Agent-Call via Claude-CLI, OpenCode (MiniMax/Ollama) oder direkter MiniMax-Text-API. Stacks sind unabhängig — fehlt ein Binary/Key, fällt nur dieser Provider aus. Bezahlte Lektionen des Vorgängers, hier eingebaut statt nachgepatcht: - OpenCode-Starts staffeln (Session-DB-Kollision: "database is locked" bei Parallelstart) - Prozessgruppen-Kill (CLI-Kinder halten sonst Pipes offen, communicate() hängt) - RAM-Gate vor CLI-Spawns (~310 MB RSS je Prozess) - Prompt via Tempdatei statt argv (ARG_MAX) - Token-Accounting auch bei Timeout/Fehler — Verschwendung bleibt sichtbar """ import asyncio import logging import os import re import shutil import signal import sqlite3 import tempfile import time import urllib.request from pathlib import Path import httpx from config import (MAX_CONCURRENT_AGENTS, MAX_CONCURRENT_API_AGENTS, OPENCODE_START_DELAY, PROVIDERS, RAM_MIN_FREE_PCT, resolve_role) log = logging.getLogger("creator2.agents") _batch_sem = asyncio.Semaphore(MAX_CONCURRENT_AGENTS) _api_sem = asyncio.Semaphore(MAX_CONCURRENT_API_AGENTS) _active: dict[str, float] = {} # key → Startzeit (Anzeige) _prozesse: dict[str, asyncio.subprocess.Process] = {} _abgebrochen: set[str] = set() # Key-Präfixe abgebrochener Läufe _pausiert: set[str] = set() # Key-Präfixe sanft pausierter Läufe # ── Globale 429-Bremse: Backoff pro Call reicht nicht — 28 Parallel-Calls # kollidieren nach der Wartezeit sofort wieder. Ein 429 drosselt deshalb ALLE: # Cooldown für neue Starts + API-Breite halbieren; Erfolge heben sie langsam. ── _cooldown_bis = 0.0 # time.monotonic(): vorher startet nichts Neues _api_breite = MAX_CONCURRENT_API_AGENTS _api_inflight = 0 _api_erfolge = 0 def drossel_melden(sekunden: float) -> None: global _cooldown_bis, _api_breite, _api_erfolge _cooldown_bis = max(_cooldown_bis, time.monotonic() + sekunden) neu = max(4, _api_breite // 2) if neu != _api_breite: log.warning("429-Bremse: API-Breite %d → %d, Cooldown %.0fs", _api_breite, neu, sekunden) _api_breite = neu _api_erfolge = 0 def _erfolg_melden() -> None: global _api_erfolge, _api_breite _api_erfolge += 1 if _api_erfolge >= 20 and _api_breite < MAX_CONCURRENT_API_AGENTS: _api_breite += 1 # langsame Erholung: 20 Erfolge kaufen +1 Breite _api_erfolge = 0 async def _drossel_warten(key: str) -> None: while True: rest = _cooldown_bis - time.monotonic() if rest <= 0 or _ist_abgebrochen(key): return await asyncio.sleep(min(rest, 1.0)) async def _api_slot(key: str) -> None: """Cooldown + adaptive Breite (unterhalb des festen Semaphors).""" global _api_inflight while True: await _drossel_warten(key) if _api_inflight < _api_breite: _api_inflight += 1 return await asyncio.sleep(0.2) class AgentErgebnis: __slots__ = ("rc", "text", "err", "tokens", "wait_s", "model") def __init__(self, rc: int, text: str, err: str, tokens: dict | None = None): self.rc, self.text, self.err, self.tokens = rc, text, err, tokens or {} self.wait_s = 0.0 # Queue-Zeit (Semaphore/Drossel) — run_agent füllt sie self.model = "" # aufgelöstes Modell — run_agent füllt es (Ledger) @property def ok(self) -> bool: return self.rc == 0 and bool(self.text.strip()) def aktive_agenten() -> list[dict]: now = time.time() return sorted([{"key": k, "laufzeit": round(now - t, 1)} for k, t in _active.items()], key=lambda a: -a["laufzeit"]) def abbrechen(prefix: str) -> None: """Alle laufenden + wartenden Agenten dieses Präfixes stoppen.""" _abgebrochen.add(prefix) for key, p in list(_prozesse.items()): if key.startswith(prefix): _kill(p) def pausieren(prefix: str) -> None: """Sanfte Pause: Wartende + Neue liefern sofort „pausiert", Laufende laufen aus und ihre Ergebnisse landen noch in der DB (Gegenstück zu abbrechen).""" _pausiert.add(prefix) def abbruch_aufheben(prefix: str) -> None: _abgebrochen.discard(prefix) _pausiert.discard(prefix) def _ist_abgebrochen(key: str) -> bool: return any(key.startswith(p) for p in _abgebrochen) def _ist_pausiert(key: str) -> bool: return any(key.startswith(p) for p in _pausiert) def provider_verfuegbar(provider: str) -> bool: cfg = PROVIDERS.get(provider) if not cfg: return False if shutil.which(cfg["cli"]) is None: return False if cfg.get("env_key") and not os.environ.get(cfg["env_key"]): return False if cfg.get("check_url"): try: urllib.request.urlopen(cfg["check_url"], timeout=1) except Exception: return False return True async def run_agent(key: str, prompt: str, timeout: int, *, provider: str, role: str = "quick", capabilities: str = "none") -> AgentErgebnis: """Ein Call. capabilities: none (nur Text) | files (Read/Write/Bash) | full (+Web). Wirft asyncio.TimeoutError bei Timeout (Infra-Behandlung macht llm.py).""" if _ist_pausiert(key): # VOR dem Fake-Shortcut — sonst ist Pause untestbar return AgentErgebnis(1, "", "pausiert") if os.getenv("CREATOR_FAKE_AGENTS"): import fake_agents return await fake_agents.antwort(key, prompt, capabilities) if _ist_abgebrochen(key): return AgentErgebnis(1, "", "abgebrochen") run_provider = provider provider, model = resolve_role(run_provider, role) if provider != run_provider and not provider_verfuegbar(provider): provider, model = run_provider, PROVIDERS[run_provider].get(role, "") if not model: return AgentErgebnis(1, "", f"kein Modell für Rolle {role} ({provider})") use_api = (capabilities == "none" and PROVIDERS[provider]["cli"] == "opencode" and model.split("/", 1)[0] in ("minimax", "minimax-kalt") and bool(os.environ.get("MINIMAX_API_KEY"))) sem = _api_sem if use_api else _batch_sem global _api_inflight warte_start = time.time() # Queue-Zeit getrennt ausweisen (Ledger wait_ms) async with sem: if _ist_pausiert(key): # wartete im Semaphor, als die Pause kam return AgentErgebnis(1, "", "pausiert") if _ist_abgebrochen(key): return AgentErgebnis(1, "", "abgebrochen") _active[key] = time.time() try: if use_api: await _api_slot(key) # belegt _api_inflight — bei Abbruch freigeben if _ist_abgebrochen(key): _api_inflight -= 1 return AgentErgebnis(1, "", "abgebrochen") wait_s = time.time() - warte_start try: res = await _text_api(key, prompt, timeout, model, role) finally: _api_inflight -= 1 if res.rc == 0: _erfolg_melden() res.wait_s = wait_s res.model = model return res await _drossel_warten(key) if _ist_abgebrochen(key): # Nachzügler nach dem Cooldown nicht mehr feuern return AgentErgebnis(1, "", "abgebrochen") wait_s = time.time() - warte_start if PROVIDERS[provider]["cli"] == "opencode": res = await _opencode(key, prompt, timeout, model, capabilities) else: res = await _claude_cli(key, prompt, timeout, model, capabilities) res.wait_s = wait_s res.model = model return res finally: _active.pop(key, None) # ── CLI-Pfade ───────────────────────────────────────────────────────────────── _CLAUDE_TOOLS = {"full": "Write,Bash,Read,WebSearch,WebFetch", "files": "Read,Bash,Write", "none": None} _OPENCODE_AGENTS = {"full": "full", "files": "files", "none": "text"} _start_lock = asyncio.Lock() _next_start = 0.0 async def _start_slot() -> None: """Token-Bucket: CLI-Starts um OPENCODE_START_DELAY spreizen, ohne globalen Konvoi.""" global _next_start loop = asyncio.get_running_loop() async with _start_lock: now = loop.time() bei = max(now, _next_start) _next_start = bei + OPENCODE_START_DELAY await asyncio.sleep(max(0.0, bei - now)) def _ram_frei_pct() -> float | None: try: text = Path("/proc/meminfo").read_text() m = dict(re.findall(r"^(MemTotal|MemAvailable):\s+(\d+)", text, re.MULTILINE)) return int(m["MemAvailable"]) * 100 / int(m["MemTotal"]) except Exception: return None async def _ram_gate(key: str) -> bool: """Unter RAM_MIN_FREE_PCT frei warten neue Spawns (fail-open ohne /proc).""" if RAM_MIN_FREE_PCT <= 0: return True while True: if _ist_abgebrochen(key): return False pct = _ram_frei_pct() if pct is None or pct >= RAM_MIN_FREE_PCT or not _prozesse: return True await asyncio.sleep(2) def _kill(process) -> None: try: os.killpg(os.getpgid(process.pid), signal.SIGKILL) except (ProcessLookupError, PermissionError): try: process.kill() except ProcessLookupError: pass async def _spawn(key: str, cmd: list[str], stdin_data: bytes | None, timeout: int, env: dict | None = None) -> AgentErgebnis: if not await _ram_gate(key): return AgentErgebnis(1, "", "abgebrochen") await _start_slot() process = await asyncio.create_subprocess_exec( *cmd, stdin=asyncio.subprocess.PIPE if stdin_data is not None else asyncio.subprocess.DEVNULL, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, start_new_session=True, env=env) track = key n = 2 while track in _prozesse: track = f"{key}~{n}" n += 1 _prozesse[track] = process start = time.monotonic() try: try: stdout, stderr = await asyncio.wait_for(process.communicate(input=stdin_data), timeout=timeout) except (asyncio.TimeoutError, asyncio.CancelledError): # Timeout ODER Cancel (Hedge-Verlierer, LaufPause) → Prozess killen, # sonst rechnet ein verwaister CLI-Prozess unsichtbar weiter _kill(process) try: await asyncio.wait_for(process.wait(), timeout=5) except (asyncio.TimeoutError, asyncio.CancelledError): pass raise log.info("agent %s: exit %s nach %.1fs", key, process.returncode, time.monotonic() - start) return AgentErgebnis(process.returncode or 0, stdout.decode("utf-8", errors="replace"), stderr.decode("utf-8", errors="replace")) finally: if _prozesse.get(track) is process: del _prozesse[track] async def _claude_cli(key: str, prompt: str, timeout: int, model: str, capabilities: str) -> AgentErgebnis: cmd = ["claude", "-p", "--model", model] if _CLAUDE_TOOLS.get(capabilities): cmd += ["--allowedTools", _CLAUDE_TOOLS[capabilities]] cmd += ["--dangerously-skip-permissions"] return await _spawn(key, cmd, prompt.encode(), timeout) _OPENCODE_DB = Path.home() / ".local" / "share" / "opencode" / "opencode.db" def _opencode_tokens(key: str) -> dict | None: """Token-Zähler der Session (run --title=key). Best effort, nie fatal.""" try: con = sqlite3.connect(f"file:{_OPENCODE_DB}?mode=ro", uri=True, timeout=1) try: row = con.execute( "SELECT tokens_input, tokens_output, tokens_cache_read, tokens_cache_write" " FROM session WHERE title=? ORDER BY time_created DESC LIMIT 1", (key,)).fetchone() finally: con.close() except Exception: return None if row is None: return None return {"input": int(row[0] or 0), "output": int(row[1] or 0), "cache_read": int(row[2] or 0), "cache_write": int(row[3] or 0)} async def _opencode(key: str, prompt: str, timeout: int, model: str, capabilities: str) -> AgentErgebnis: with tempfile.NamedTemporaryFile("w", suffix=".md", delete=False, encoding="utf-8") as f: f.write(prompt) pfad = Path(f.name) cmd = ["opencode", "run", "Folge exakt den Anweisungen in der angehängten Datei. Sie sind der vollständige Auftrag.", "-m", model, "--agent", _OPENCODE_AGENTS.get(capabilities, "text"), "--dangerously-skip-permissions", "--title", key, "-f", str(pfad)] try: try: res = await _spawn(key, cmd, None, timeout) except asyncio.TimeoutError: # Timeout heißt nicht 0 Tokens — die Session hat oft schon zigtausend # erzeugt. Als infra (err="timeout") mit echten Tokens zurückgeben, # damit Budget/Ledger stimmen. CancelledError NICHT fangen (Hedge). tokens = await asyncio.to_thread(_opencode_tokens, key) return AgentErgebnis(1, "", "timeout", tokens or {}) res.text = _clean_opencode(res.text) res.tokens = await asyncio.to_thread(_opencode_tokens, key) or {} return res finally: pfad.unlink(missing_ok=True) _ANSI = re.compile(r"\x1b\[[0-9;]*m") def _clean_opencode(text: str) -> str: text = _ANSI.sub("", text) lines = text.splitlines() while lines and (not lines[0].strip() or lines[0].lstrip().startswith(">")): lines.pop(0) return "\n".join(lines).strip() # ── Direkte MiniMax-Text-API (spart den ~310-MB-Prozess für tool-lose Calls) ── _API_URL = "https://api.minimax.io/anthropic/v1/messages" _API_MAX_TOKENS = 32_000 _API_MODEL_OPTS = { "minimax-kalt/MiniMax-M3": {"temperature": 0.2, "thinking": {"type": "disabled"}}, # thinking AUS: ungebremst dachte das Modell bei großen Extraktions-Chunks # 8k–32k Output-Tokens lang (aak Lauf 15: 4 Calls liefen ins 32k-Cap, leere # Antwort nach bis zu 44 min — Extraktion war 53 von 63 min der Ebene). "minimax-kalt/MiniMax-M2.7-highspeed": {"temperature": 0.3, "thinking": {"type": "disabled"}}, # native Route (judge/guide): ungedrosselt liefen Merge-Judges ins Output-Cap # (58–71 % Parse-Fehler). role=guide (Writer) behält Thinking — _text_api # nimmt die thinking-Option dort wieder raus. "minimax/MiniMax-M3": {"temperature": 0.2, "thinking": {"type": "disabled"}}, } async def _text_api(key: str, prompt: str, timeout: int, model: str, role: str = "judge") -> AgentErgebnis: opts = dict(_API_MODEL_OPTS.get(model, {})) if role == "guide": opts.pop("thinking", None) # Writer braucht Denkraum — nur Judges kalt body = {"model": model.split("/", 1)[1], "max_tokens": _API_MAX_TOKENS, "messages": [{"role": "user", "content": prompt}], **opts} headers = {"x-api-key": os.environ.get("MINIMAX_API_KEY", ""), "anthropic-version": "2023-06-01"} try: async with httpx.AsyncClient(timeout=httpx.Timeout(timeout, connect=30)) as client: resp = await asyncio.wait_for(client.post(_API_URL, json=body, headers=headers), timeout=timeout) except httpx.TimeoutException: raise asyncio.TimeoutError except httpx.HTTPError as e: return AgentErgebnis(1, "", f"{type(e).__name__}: {e}") if resp.status_code == 429: try: retry_after = float(resp.headers.get("retry-after", "15")) except ValueError: retry_after = 15.0 drossel_melden(min(retry_after, 60.0)) return AgentErgebnis(1, "", "HTTP 429: rate limited") if resp.status_code != 200: return AgentErgebnis(1, "", f"HTTP {resp.status_code}: {resp.text[:300]}") data = resp.json() text = "".join(b.get("text", "") for b in data.get("content", []) if b.get("type") == "text") u = data.get("usage") or {} tokens = {"input": int(u.get("input_tokens") or 0), "output": int(u.get("output_tokens") or 0), "cache_read": int(u.get("cache_read_input_tokens") or 0), "cache_write": int(u.get("cache_creation_input_tokens") or 0)} if not text.strip(): return AgentErgebnis(1, "", f"leere Antwort (stop={data.get('stop_reason')})", tokens) return AgentErgebnis(0, text, "", tokens)