Files
mission-control-v2/deploy/governor/governor.py
T
Hitonabi 798de0be8f Governor Phase 0 fertig + Phase 2: Hart-Deckel default + Sprach-Signal an Lucy
Phase 0 abgeschlossen:
- Hart-Deckel jetzt Default AN (Soft+5000), da der weiche Schnitt allein bei
  Weiterarbeit ueber die Grenze eine Fassade erzeugt (P0-Befund). GOV_HARD_CEILING=0
  schaltet ihn aus. gov-ctl reicht Arg 2 nur bei Bedarf durch.
- README: Empfehlung/Doku auf Default-an aktualisiert, Commit-Msg-Restpunkt notiert
  (--no-auto-commits als saubere Option).

Phase 2 (Voice-Hook): beim Feuern (soft/hard) POSTet der Governor best-effort +
gedrosselt (Default 300 s) eine Meldung an Lucys vorhandene Announce-Pipeline
(POST :9001/api/voice/announce, source=governor, priority=normal). Lucy pollt,
dedupliziert und spricht sie via lokales TTS -- KEIN Lucy-Code noetig. E2E bewiesen:
Feuern -> ANNOUNCE status=200 -> Eintrag id 394 source=governor in der Queue.
GOV_ANNOUNCE_URL="" schaltet es ab. announce-test.sh beigelegt.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-24 20:06:10 +02:00

421 lines
17 KiB
Python

#!/usr/bin/env python3
"""Governor — duenner, zustandsloser Token-Waechter-Proxy (Phase 0).
Sitzt zwischen einem Coding-Agenten (Aider) und dem Modell-Endpoint (llama-swap
:8080). Reicht ALLES unveraendert durch — mit einer Ausnahme bei
/v1/chat/completions: er schaetzt die Token-Groesse der Anfrage (= Sessiongroesse,
weil die ganze Historie jede Runde mitkommt) und handelt nach zwei Schwellen:
est >= SCHWELLE (soft): haengt eine Stopp-Anweisung als letzte User-Nachricht an
("SAVEPOINT.md finalisieren + stoppen") und leitet weiter. Das Modell schreibt
EINEN ehrlichen Abschluss-Savepoint. Loggt FIRED.
est >= HART-DECKEL (optional): antwortet SELBST mit einer kurzen Stopp-Nachricht,
OHNE das Modell zu fragen. Verhindert, dass ueber die Grenze hinaus
weitergearbeitet wird — genau das erzeugte in Tests eine Fassade. Loggt HARDSTOP.
Bewusst nur Standardbibliothek: kein pip, kein venv, laeuft mit System-python3.
Bewusst zustandslos: jede Anfrage wird fuer sich bewertet; keine Sitzungs-DB.
Konfiguration per Umgebungsvariablen (alle optional):
GOV_PORT Listen-Port (Default 8100)
GOV_HOST Listen-Adresse (Default 0.0.0.0)
GOV_UPSTREAM Modell-Endpoint (Default http://127.0.0.1:8080)
GOV_THRESHOLD Soft-Schwelle fuer den Einschub (Default 25000)
GOV_HARD_CEILING Hart-Deckel; 0 = aus (Default: Soft+5000, AN)
GOV_CHARS_PER_TOKEN Heuristik Zeichen->Token (Default 3.5, kalibriert)
GOV_LOG Logdatei (zusaetzlich zu stdout) (Default ./governor.log)
GOV_DIRECTIVE Text des Soft-Einschubs (sonst Default unten)
GOV_HARDSTOP_MSG Text der Hart-Stopp-Antwort (sonst Default unten)
GOV_ANNOUNCE_URL Lucy-Sprach-Signal-Endpunkt; "" = aus (Default :9001/api/voice/announce)
GOV_ANNOUNCE_THROTTLE Sekunden zwischen Signalen (Default 300)
GOV_ANNOUNCE_TEXT Text des Sprach-Signals (sonst Default)
"""
import http.client
import json
import os
import re
import threading
import time
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")
UPSTREAM = os.environ.get("GOV_UPSTREAM", "http://127.0.0.1:8080")
THRESHOLD = int(os.environ.get("GOV_THRESHOLD", "25000"))
# Hart-Deckel: Default AN (Soft + 5000 = ein Finalisier-Zug Luft), weil der weiche
# Schnitt allein bei Weiterarbeit ueber die Grenze eine Fassade erzeugt (Befund P0).
# Explizit setzbar; GOV_HARD_CEILING=0 schaltet ihn aus.
_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.5"))
LOG_PATH = os.environ.get("GOV_LOG", os.path.join(os.getcwd(), "governor.log"))
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)
# Sprach-Signal an Lucy (Phase 2): beim Feuern POSTet der Governor eine Meldung an die
# vorhandene MC2-Announce-Pipeline (:9001). Lucy pollt sie ohnehin, dedupliziert und
# spricht sie (gated durch ihren "Box-Meldungen laut"-Schalter). Best-effort, gedrosselt
# gegen die Pro-Runde-Feuerung. GOV_ANNOUNCE_URL="" schaltet das Signal ab.
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")) # Sekunden
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 # Zeitstempel der letzten Meldung (Drossel)
_an = urlparse(ANNOUNCE_URL) if ANNOUNCE_URL else None
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 estimate_tokens(messages) -> int:
"""Grobe, aber stabile Heuristik: Zeichen aller Nachrichteninhalte / CPT.
Gegen die echten prompt_tokens aus der Antwort kalibriert (CPT=3.5 traf am
24.07. auf ~1-3 % genau). Zaehlt Text in String- und Multimodal-Listen-Inhalten;
kleiner Aufschlag je Nachricht fuer Rollen-/Template-Overhead.
"""
chars = 0
for m in messages or []:
chars += 4
content = m.get("content") if isinstance(m, dict) else None
if isinstance(content, str):
chars += len(content)
elif isinstance(content, list):
for part in content:
if isinstance(part, dict) and isinstance(part.get("text"), str):
chars += len(part["text"])
return int(chars / CHARS_PER_TOKEN)
def _post_announce(est) -> None:
"""POSTet die Meldung an die MC2-Announce-Pipeline. Laeuft im Hintergrund-Thread."""
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: # noqa: BLE001
log(f"ANNOUNCE fehlgeschlagen: {exc!r}")
def maybe_announce(est) -> None:
"""Sprach-Signal an Lucy ausloesen — gedrosselt, damit die Pro-Runde-Feuerung
nicht spammt (eine Aeusserung je Ueberschreitungs-Episode). Best-effort."""
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()
class Handler(BaseHTTPRequestHandler):
protocol_version = "HTTP/1.1"
server_version = "Governor/0.2"
def log_message(self, *args):
pass
def do_GET(self):
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 _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
# Query-String vor der Endpunkt-Erkennung abschneiden (sonst umgeht
# z. B. ?api-version=... den Governor).
clean_path = path.split("?", 1)[0]
is_chat = clean_path.rstrip("/").endswith("/chat/completions")
action = "passthrough"
est = None
streaming = False
model = ""
if is_chat and body:
action, body, est, streaming, model = self._decide(body)
if action in ("soft", "hard"):
maybe_announce(est) # Sprach-Signal an Lucy (gedrosselt)
# HARTER STOPP: selbst antworten, Upstream nie fragen.
if action == "hard":
self._send_canned_stop(model, streaming, est)
log(f"chat est={est} thr={THRESHOLD} hard={HARD_CEILING} HARDSTOP "
f"stream={streaming} status=200")
return
# Header fuer Upstream aufbereiten.
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=600)
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()
# hop-by-hop + Laenge raus; Date/Server setzt send_response schon selbst.
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 bis 64 KB oder Stream-Ende puffern
# und streamendes Aider die ganze Generierung 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:
exact_s = str(exact) if exact is not None else "-"
flag = "FIRED" if action == "soft" else "ok"
log(f"chat est={est} exact={exact_s} thr={THRESHOLD} {flag} "
f"stream={streaming} status={resp.status}")
def _decide(self, body: bytes):
"""Aktion bestimmen: passthrough | soft (Einschub) | hard (Selbstantwort).
Rueckgabe: (action, body, est, streaming, model).
"""
try:
data = json.loads(body)
except (ValueError, UnicodeDecodeError):
return "passthrough", body, None, False, ""
if not isinstance(data, dict):
return "passthrough", body, None, False, ""
messages = data.get("messages")
streaming = bool(data.get("stream"))
model = data.get("model", "") or ""
est = estimate_tokens(messages if isinstance(messages, list) else [])
if HARD_CEILING > 0 and est >= HARD_CEILING:
return "hard", body, est, streaming, model
if est >= THRESHOLD and isinstance(messages, list):
# Sichere Substitution statt str.format: ein operator-gesetzter
# GOV_DIRECTIVE mit { } (JSON/Code-Beispiel) darf nicht crashen.
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
return "passthrough", body, est, streaming, model
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:
try:
data = json.dumps({"error": msg}).encode()
self.send_response(status)
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)
except OSError:
pass
@staticmethod
def _scan_prompt_tokens(tail: bytearray):
if not tail:
return None
try:
text = tail.decode("utf-8", errors="ignore")
except Exception: # noqa: BLE001
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 startet auf {HOST}:{PORT} -> {UPSTREAM} | Soft={THRESHOLD} "
f"Hart={hard} | CPT={CHARS_PER_TOKEN} | 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())