#!/usr/bin/env python3 """Governor v2 — Token-Waechter-Proxy vor dem MC2-Gateway. Sitzt zwischen den Coding-Agenten (OpenCode/Zed, Nacht-Laeufe, Hermes-Worker) und dem MC2-Gateway (:9001). Reicht ALLES unveraendert durch — mit einer Ausnahme bei /v1/chat/completions: er bestimmt die Groesse der Anfrage (= Sitzungsgroesse, weil die ganze Historie jede Runde mitkommt) und handelt nach zwei Schwellen: est >= SOFT: haengt eine Stopp-Anweisung als letzte User-Nachricht an ("SAVEPOINT.md finalisieren + stoppen") und leitet weiter. Loggt FIRED. est >= HART: antwortet SELBST mit einer kurzen Stopp-Nachricht, OHNE das Modell zu fragen. Verhindert Fassaden jenseits der Grenze. Loggt HARDSTOP. --- Was v2 gegenueber v0.2 aendert (25.07.2026) --------------------------------------- 1. EHRLICH ZAEHLEN. v0.2 zaehlte nur Text in `messages` und ignorierte `tools`/ `tool_calls`. Bei werkzeugdichten Agenten lag es um Faktor 3 daneben (gemessen im eigenen Log: est=3353 exact=10224). v2 zaehlt den GANZEN Anfragekoerper — inklusive Werkzeug-Schemata, Werkzeug-Aufrufe und Werkzeug-Ergebnisse. 2. SELBST-KALIBRIERUNG. Aus jeder Antwort liest der Governor die echten `usage.prompt_tokens` und korrigiert damit sein Zeichen-pro-Token-Verhaeltnis — pro Modell, gleitend. Die Schaetzung wird also im Betrieb immer genauer, statt auf einem einmal geratenen Wert festzuhaengen. 3. TOOL-CALL-SICHERER EINSCHUB. Der Soft-Einschub wird NUR angehaengt, wenn die Nachrichtenkette das erlaubt (letzte Nachricht ist nicht ein Assistant mit offenen tool_calls und keine tool-Antwort). Sonst wartet er auf die naechste Runde. Ohne diese Pruefung zerbricht der Einschub bei OpenCode die Werkzeug-Reihenfolge. 4. STATUS-ENDPUNKT. GET /governor/status liefert Zaehlerstand, Kalibrierung und die letzten Laeufe als JSON — Datenquelle fuer die MC2-Oberflaeche, das OpenCode-Plugin und Lucys `loop_status`. Bewusst nur Standardbibliothek: kein pip, kein venv, laeuft mit System-python3. Bewusst ohne Datenbank: ein kleiner Ring im Speicher, mehr braucht es nicht. Konfiguration per Umgebungsvariablen (alle optional): GOV_PORT Listen-Port (Default 8100) GOV_HOST Listen-Adresse (Default 0.0.0.0) GOV_UPSTREAM Ziel (Default http://127.0.0.1:9001) GOV_THRESHOLD Soft-Schwelle fuer den Einschub (Default 45000) GOV_HARD_CEILING Hart-Deckel; 0 = aus (Default: Soft+5000, AN) GOV_CHARS_PER_TOKEN Startwert Zeichen->Token (Default 3.2, danach gelernt) GOV_CALIBRATE Selbst-Kalibrierung an/aus (Default 1) GOV_LOG Logdatei (zusaetzlich zu stdout) (Default ./governor.log) GOV_DIRECTIVE Text des Soft-Einschubs GOV_HARDSTOP_MSG Text der Hart-Stopp-Antwort GOV_ANNOUNCE_URL Lucy-Sprach-Signal; "" = aus (Default :9001/api/voice/announce) GOV_ANNOUNCE_THROTTLE Sekunden zwischen Signalen (Default 300) GOV_ANNOUNCE_TEXT Text des Sprach-Signals GOV_EXEMPT_MODELS Modelle ohne Schnitt, kommasepariert (Default: hermes,fast,embed, reranker,vision,scout — Lucys Alltag wird nie unterbrochen) """ import http.client import json import os import re import threading import time from collections import deque from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from urllib.parse import urlparse # ---- Konfiguration --------------------------------------------------------- PORT = int(os.environ.get("GOV_PORT", "8100")) HOST = os.environ.get("GOV_HOST", "0.0.0.0") # Ziel ist das MC2-Gateway, NICHT llama-swap direkt: so bleiben MC2s Rollen-Aliase, # Bild-Weiche und Telemetrie erhalten. Der Governor ist eine Schicht davor, kein Ersatz. UPSTREAM = os.environ.get("GOV_UPSTREAM", "http://127.0.0.1:9001") THRESHOLD = int(os.environ.get("GOV_THRESHOLD", "45000")) _hard_env = os.environ.get("GOV_HARD_CEILING") HARD_CEILING = (THRESHOLD + 5000) if _hard_env is None else int(_hard_env) CHARS_PER_TOKEN = float(os.environ.get("GOV_CHARS_PER_TOKEN", "3.2")) CALIBRATE = os.environ.get("GOV_CALIBRATE", "1") != "0" LOG_PATH = os.environ.get("GOV_LOG", os.path.join(os.getcwd(), "governor.log")) # Lucys Alltagsmodelle bekommen NIE einen Savepoint-Einschub: sie fuehren Gespraeche, # keine Bau-Sitzungen. Nur die Coding-Rollen laufen gegen die Schwelle. _DEFAULT_EXEMPT = "hermes,fast,embed,reranker,vision,scout" EXEMPT_MODELS = {m.strip().lower() for m in os.environ.get("GOV_EXEMPT_MODELS", _DEFAULT_EXEMPT).split(",") if m.strip()} DEFAULT_DIRECTIVE = ( "[GOVERNOR — SITZUNGS-LIMIT ERREICHT] Der Kontext dieser Sitzung ist auf ~{est} " "Tokens gewachsen (Limit {threshold}). Beginne oder setze JETZT KEINE weiteren " "Code-Aenderungen fort. Stattdessen, in dieser Reihenfolge:\n" "1. Aktualisiere SAVEPOINT.md so, dass es den aktuellen Stand vollstaendig festhaelt: " "was WIRKLICH erledigt ist (nur was im Code steht — nichts aus Absicht oder git-" "Nachrichten ableiten), der genaue naechste Schritt, offene Fragen und alle " "Stolpersteine — genug, dass eine frische Sitzung ohne jede Erinnerung allein aus " "SAVEPOINT.md plus git-Historie sauber weitermachen kann.\n" "2. Halte dann an und sage dem Nutzer in einem Satz, dass er eine frische Sitzung " "starten soll. Gib ausser der SAVEPOINT.md-Aktualisierung und diesem Hinweis nichts aus." ) DIRECTIVE = os.environ.get("GOV_DIRECTIVE", DEFAULT_DIRECTIVE) DEFAULT_HARDSTOP = ( "[GOVERNOR — HARTER STOPP] Das Sitzungs-Limit ist ueberschritten und der Savepoint " "sollte bereits finalisiert sein. Diese Sitzung nimmt keine weiteren Auftraege mehr an. " "Bitte starte eine FRISCHE Sitzung — sie liest SAVEPOINT.md und die git-Historie und " "macht sauber weiter. (Keine Code-Aenderung in dieser Antwort.)" ) HARDSTOP_MSG = os.environ.get("GOV_HARDSTOP_MSG", DEFAULT_HARDSTOP) ANNOUNCE_URL = os.environ.get("GOV_ANNOUNCE_URL", "http://127.0.0.1:9001/api/voice/announce") ANNOUNCE_THROTTLE = float(os.environ.get("GOV_ANNOUNCE_THROTTLE", "300")) DEFAULT_ANNOUNCE = ( "Commander, die Coding-Sitzung wird voll — ungefähr {est} Tokens. Ich sichere den " "Stand im Savepoint; am besten fangen wir gleich frisch an." ) ANNOUNCE_TEXT = os.environ.get("GOV_ANNOUNCE_TEXT", DEFAULT_ANNOUNCE) up = urlparse(UPSTREAM) UP_HOST = up.hostname or "127.0.0.1" UP_PORT = up.port or 80 HOP_BY_HOP = { "connection", "keep-alive", "proxy-authenticate", "proxy-authorization", "te", "trailers", "transfer-encoding", "upgrade", } _log_lock = threading.Lock() _PROMPT_TOKENS_RE = re.compile(r'"prompt_tokens"\s*:\s*(\d+)') _announce_lock = threading.Lock() _last_announce = 0.0 _an = urlparse(ANNOUNCE_URL) if ANNOUNCE_URL else None # ---- Zustand (klein, im Speicher) ------------------------------------------ # Kalibrierung je Modell: gleitender Mittelwert von zeichen/echte_tokens. Startwert ist # GOV_CHARS_PER_TOKEN; jede Antwort mit usage zieht ihn Richtung Wahrheit. _state_lock = threading.Lock() _cpt: dict = {} # modell -> gelerntes Zeichen-pro-Token _cpt_n: dict = {} # modell -> Anzahl Messungen _recent: deque = deque(maxlen=50) # letzte Laeufe fuer /governor/status _counters = {"chat": 0, "soft": 0, "hard": 0, "passthrough": 0, "tokens_prompt": 0, "tokens_completion": 0, "started": time.time()} CPT_MIN, CPT_MAX = 0.8, 8.0 # Schutz gegen Ausreisser def log(line: str) -> None: """Eine Zeile nach stdout UND in die Logdatei (thread-sicher).""" stamp = time.strftime("%Y-%m-%dT%H:%M:%S") msg = f"{stamp} {line}" with _log_lock: print(msg, flush=True) try: with open(LOG_PATH, "a", encoding="utf-8") as fh: fh.write(msg + "\n") except OSError: pass def cpt_for(model: str) -> float: """Aktuelles Zeichen-pro-Token-Verhaeltnis fuer ein Modell (gelernt oder Startwert).""" with _state_lock: return _cpt.get(model, CHARS_PER_TOKEN) def calibrate(model: str, chars: int, exact: int) -> None: """Aus einer echten Antwort lernen. Gleitender Mittelwert mit sanftem Gewicht — ein einzelner Ausreisser (z. B. ein riesiges Bild) verbiegt nichts.""" if not CALIBRATE or not exact or exact <= 0 or chars <= 0: return ratio = chars / exact if not (CPT_MIN <= ratio <= CPT_MAX): return with _state_lock: n = _cpt_n.get(model, 0) old = _cpt.get(model, CHARS_PER_TOKEN) # Gewicht faellt mit der Anzahl Messungen: schnell einschwingen, dann stabil. w = max(0.08, 1.0 / (n + 2)) _cpt[model] = old * (1 - w) + ratio * w _cpt_n[model] = n + 1 def body_chars(data: dict) -> int: """Zeichen des GESAMTEN Anfragekoerpers — der Kern der Ehrlichkeit. v0.2 zaehlte nur Text in `messages` und lag bei werkzeugdichten Agenten um Faktor 3 daneben, weil Werkzeug-Schemata (`tools`), Werkzeug-Aufrufe (`tool_calls`) und Werkzeug-Ergebnisse mitgeschickt werden und im Kontext genauso Platz fressen. Wir serialisieren einfach alles, was ans Modell geht. """ payload = {k: v for k, v in data.items() if k in ("messages", "tools", "tool_choice", "system", "functions")} try: return len(json.dumps(payload, ensure_ascii=False)) except (TypeError, ValueError): # Fallback: nur Nachrichtentext (nie schlechter als v0.2) chars = 0 for m in data.get("messages") or []: c = m.get("content") if isinstance(m, dict) else None if isinstance(c, str): chars += len(c) + 4 elif isinstance(c, list): for p in c: if isinstance(p, dict) and isinstance(p.get("text"), str): chars += len(p["text"]) return chars def safe_to_append(messages) -> bool: """Darf der Soft-Einschub JETZT als user-Nachricht ans Ende? Nein, wenn die Kette gerade mitten in einem Werkzeug-Austausch steckt: nach einem Assistant mit offenen `tool_calls` MUSS eine `tool`-Antwort folgen — schiebt man da eine user-Nachricht dazwischen, lehnt das Modell (bzw. das Template) die Anfrage ab oder halluziniert. Dann warten wir einfach auf die naechste Runde; die Schwelle ist ohnehin ueberschritten, es kommt in Sekunden ein neuer Zug. """ if not isinstance(messages, list) or not messages: return False last = messages[-1] if not isinstance(last, dict): return False role = last.get("role") if role == "tool": return False return not (role == "assistant" and last.get("tool_calls")) def _post_announce(est) -> None: """POSTet die Meldung an die MC2-Announce-Pipeline (Lucy spricht sie).""" try: text = (ANNOUNCE_TEXT.replace("{est}", str(est)) .replace("{threshold}", str(THRESHOLD))) body = json.dumps({"text": text, "subject": "[Governor]", "source": "governor", "priority": "normal"}).encode("utf-8") conn = http.client.HTTPConnection(_an.hostname or "127.0.0.1", _an.port or 80, timeout=4) conn.request("POST", _an.path or "/api/voice/announce", body=body, headers={"Content-Type": "application/json", "Content-Length": str(len(body))}) resp = conn.getresponse() resp.read() conn.close() log(f"ANNOUNCE -> Lucy status={resp.status} est={est}") except Exception as exc: log(f"ANNOUNCE fehlgeschlagen: {exc!r}") def maybe_announce(est) -> None: """Sprach-Signal an Lucy — gedrosselt (eine Aeusserung je Episode).""" if not _an: return global _last_announce now = time.time() with _announce_lock: if now - _last_announce < ANNOUNCE_THROTTLE: return _last_announce = now threading.Thread(target=_post_announce, args=(est,), daemon=True).start() def status_payload() -> dict: """Momentaufnahme fuer /governor/status (MC2-Oberflaeche, Plugin, Lucy).""" with _state_lock: return { "ok": True, "upstream": UPSTREAM, "soft": THRESHOLD, "hard": HARD_CEILING if HARD_CEILING > 0 else None, "uptime_s": int(time.time() - _counters["started"]), "counters": {k: v for k, v in _counters.items() if k != "started"}, "calibration": {m: {"chars_per_token": round(v, 3), "samples": _cpt_n.get(m, 0)} for m, v in _cpt.items()}, "calibration_default": CHARS_PER_TOKEN, "exempt_models": sorted(EXEMPT_MODELS), "recent": list(_recent), } class Handler(BaseHTTPRequestHandler): protocol_version = "HTTP/1.1" server_version = "Governor/2.0" def log_message(self, *args): pass def do_GET(self): if self.path.split("?", 1)[0].rstrip("/") in ("/governor/status", "/governor"): self._send_json(200, status_payload()) return self._proxy() def do_POST(self): self._proxy() def do_PUT(self): self._proxy() def do_DELETE(self): self._proxy() def do_OPTIONS(self): self._proxy() # -- Kern --------------------------------------------------------------- def _send_json(self, status: int, obj) -> None: try: data = json.dumps(obj).encode("utf-8") self.send_response(status) self.send_header("Content-Type", "application/json") self.send_header("Access-Control-Allow-Origin", "*") self.send_header("Content-Length", str(len(data))) self.send_header("Connection", "close") self.end_headers() self.wfile.write(data) except OSError: pass def _read_body(self) -> bytes: length = self.headers.get("Content-Length") if length is None: return b"" try: return self.rfile.read(int(length)) except (ValueError, OSError): return b"" def _proxy(self) -> None: body = self._read_body() path = self.path clean_path = path.split("?", 1)[0] is_chat = clean_path.rstrip("/").endswith("/chat/completions") action = "passthrough" est = None streaming = False model = "" chars = 0 if is_chat and body: action, body, est, streaming, model, chars = self._decide(body) if action in ("soft", "hard"): maybe_announce(est) if action == "hard": self._send_canned_stop(model, streaming, est) with _state_lock: _counters["chat"] += 1 _counters["hard"] += 1 _recent.appendleft({"t": int(time.time()), "model": model, "est": est, "exact": None, "action": "hard"}) log(f"chat model={model} est={est} thr={THRESHOLD} hard={HARD_CEILING} " f"HARDSTOP stream={streaming} status=200") return out_headers = {} for k, v in self.headers.items(): kl = k.lower() if kl in HOP_BY_HOP or kl in ("host", "content-length", "accept-encoding"): continue out_headers[k] = v out_headers["Host"] = f"{UP_HOST}:{UP_PORT}" out_headers["Accept-Encoding"] = "identity" if body: out_headers["Content-Length"] = str(len(body)) out_headers["Connection"] = "close" conn = None try: conn = http.client.HTTPConnection(UP_HOST, UP_PORT, timeout=900) conn.request(self.command, path, body=body or None, headers=out_headers) resp = conn.getresponse() except (OSError, http.client.HTTPException) as exc: log(f"ERROR upstream {self.command} {path}: {exc!r}") if conn is not None: conn.close() self._safe_error(502, f"governor upstream: {exc}") return self.send_response(resp.status) for k, v in resp.getheaders(): kl = k.lower() if kl in HOP_BY_HOP or kl in ("content-length", "date", "server"): continue self.send_header(k, v) self.send_header("Connection", "close") self.end_headers() tail = bytearray() try: while True: # read1() gibt jedes Upstream-Stueck sofort zurueck (echtes SSE- # Durchreichen); read() wuerde puffern und Streaming haengen lassen. chunk = resp.read1(65536) if not chunk: break self.wfile.write(chunk) self.wfile.flush() tail.extend(chunk) if len(tail) > 16384: del tail[:-16384] except OSError: pass finally: conn.close() exact = self._scan_prompt_tokens(tail) if is_chat: if exact: calibrate(model, chars, exact) with _state_lock: _counters["chat"] += 1 _counters["soft" if action == "soft" else "passthrough"] += 1 if exact: _counters["tokens_prompt"] += exact _recent.appendleft({"t": int(time.time()), "model": model, "est": est, "exact": exact, "action": action}) exact_s = str(exact) if exact is not None else "-" flag = "FIRED" if action == "soft" else ("skip" if action == "defer" else "ok") log(f"chat model={model} est={est} exact={exact_s} cpt={cpt_for(model):.2f} " f"thr={THRESHOLD} {flag} stream={streaming} status={resp.status}") def _decide(self, body: bytes): """Aktion bestimmen. Rueckgabe: (action, body, est, streaming, model, chars).""" try: data = json.loads(body) except (ValueError, UnicodeDecodeError): return "passthrough", body, None, False, "", 0 if not isinstance(data, dict): return "passthrough", body, None, False, "", 0 messages = data.get("messages") streaming = bool(data.get("stream")) model = (data.get("model") or "").strip() chars = body_chars(data) est = int(chars / max(cpt_for(model), 0.1)) # Lucys Alltagsmodelle laufen nie gegen die Schwelle — ein Gespraech ist keine # Bau-Sitzung. Wir zaehlen sie trotzdem mit (Kalibrierung + Telemetrie). base = model.split("/")[-1].lower() if base in EXEMPT_MODELS: return "passthrough", body, est, streaming, model, chars if not isinstance(messages, list): return "passthrough", body, est, streaming, model, chars has_warned = False marker = "[GOVERNOR — SITZUNGS-LIMIT ERREICHT]" for m in reversed(messages): if m.get("role") == "user" and isinstance(m.get("content"), str) and marker in m["content"]: has_warned = True break is_safe = safe_to_append(messages) # 1. Absolutes Not-Aus bei komplettem Amoklauf if HARD_CEILING > 0 and est >= HARD_CEILING + 10000: return "hard", body, est, streaming, model, chars # 2. Sind wir am Limit? if est >= THRESHOLD: if not is_safe: # Agent arbeitet gerade an einer Tool-Kette. Auf keinen Fall abbrechen! return "defer", body, est, streaming, model, chars # Es ist sicher (Tool-Kette beendet oder Agent wartet auf Input). if HARD_CEILING > 0 and est >= HARD_CEILING: if has_warned: # Wir haben ihn schon gewarnt, er macht trotzdem weiter -> harter Schnitt. return "hard", body, est, streaming, model, chars else: # Er hat wegen einer langen Tool-Kette direkt das Hard-Limit ueberschritten. # Gib ihm trotzdem noch den einen Finalisier-Zug (Soft). pass # Soft-Limit greift if not has_warned: directive = (DIRECTIVE.replace("{est}", str(est)) .replace("{threshold}", str(THRESHOLD))) messages.append({"role": "user", "content": directive}) data["messages"] = messages return "soft", json.dumps(data).encode("utf-8"), est, streaming, model, chars return "passthrough", body, est, streaming, model, chars def _send_canned_stop(self, model: str, streaming: bool, est) -> None: """OpenAI-kompatible Stopp-Antwort selbst erzeugen (kein Upstream-Call).""" created = int(time.time()) usage = {"prompt_tokens": est or 0, "completion_tokens": 0, "total_tokens": est or 0} try: if streaming: self.send_response(200) self.send_header("Content-Type", "text/event-stream") self.send_header("Cache-Control", "no-cache") self.send_header("Connection", "close") self.end_headers() def sse(obj): self.wfile.write(b"data: " + json.dumps(obj).encode() + b"\n\n") self.wfile.flush() base = {"id": "governor-hardstop", "object": "chat.completion.chunk", "created": created, "model": model} sse({**base, "choices": [{"index": 0, "delta": {"role": "assistant"}, "finish_reason": None}]}) sse({**base, "choices": [{"index": 0, "delta": {"content": HARDSTOP_MSG}, "finish_reason": None}]}) sse({**base, "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}], "usage": usage}) self.wfile.write(b"data: [DONE]\n\n") self.wfile.flush() else: payload = { "id": "governor-hardstop", "object": "chat.completion", "created": created, "model": model, "choices": [{"index": 0, "finish_reason": "stop", "message": {"role": "assistant", "content": HARDSTOP_MSG}}], "usage": usage, } data = json.dumps(payload).encode("utf-8") self.send_response(200) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(data))) self.send_header("Connection", "close") self.end_headers() self.wfile.write(data) self.wfile.flush() except OSError: pass def _safe_error(self, status: int, msg: str) -> None: self._send_json(status, {"error": msg}) @staticmethod def _scan_prompt_tokens(tail: bytearray): if not tail: return None try: text = tail.decode("utf-8", errors="ignore") except Exception: return None matches = _PROMPT_TOKENS_RE.findall(text) if not matches: return None try: return int(matches[-1]) except ValueError: return None def main() -> int: log_dir = os.path.dirname(LOG_PATH) if log_dir and not os.path.isdir(log_dir): try: os.makedirs(log_dir, exist_ok=True) except OSError: pass server = ThreadingHTTPServer((HOST, PORT), Handler) server.daemon_threads = True hard = HARD_CEILING if HARD_CEILING > 0 else "aus" log(f"Governor v2 startet auf {HOST}:{PORT} -> {UPSTREAM} | Soft={THRESHOLD} " f"Hart={hard} | CPT-Start={CHARS_PER_TOKEN} kalibrierend={CALIBRATE} | " f"ausgenommen={sorted(EXEMPT_MODELS)} | Log={LOG_PATH}") try: server.serve_forever() except KeyboardInterrupt: log("Governor beendet (SIGINT).") finally: server.server_close() return 0 if __name__ == "__main__": raise SystemExit(main())