Files
mission-control-v2/deploy/governor/governor.py
T
Hitonabi 7fac17ed9a Umbau Stufe 1-5: Governor v2 + OpenCode-Plugin + Subagenten-Mannschaft
Die Intelligenz sitzt nicht mehr NEBEN dem Coden, sondern DRIN. Alles auf der
Box gemessen, nicht angenommen.

Stufe 1 — Speicher-Haushalt (llama-swap-Config, nicht im Repo):
  Coder dauerwarm statt ttl 600 (Kaltstart 23 s -> 491 ms), Devstral als
  'kritiker' verdrahtet, Swap 11 GB -> 0. Zwei Befunde: Kontext 131k->65k spart
  bei Qwen3-Next 0 GB (Hybrid-Attention), und llama-swap haelt nur EINE Gruppe
  resident -> alles Ko-Residente in dieselbe Gruppe. `persistent: true` verhindert
  dabei das Freiraeumen vor grossen Modellen -> ausgeloester Kernel-OOM, behoben
  durch persistent:false + TTLs (Vision laedt jetzt in 10 s statt 117 s + Absturz).

Stufe 2 — Governor v2 (deploy/governor/):
  Steht jetzt IM Pfad (:8100 -> MC2 :9001) statt daneben. Zaehlt den GANZEN
  Anfragekoerper inkl. tools/tool_calls und kalibriert sich aus den echten
  usage.prompt_tokens jeder Antwort nach: Schaetzfehler 200 % -> 0,3 %.
  Soft-Einschub nur noch, wenn die Nachrichtenkette es erlaubt (kein Dazwischen-
  funken in offene tool_calls). Neuer Status-Endpunkt + systemd-Unit.
  Lucys Alltagsmodelle sind vom Schnitt ausgenommen.

Stufe 3 — OpenCode-Plugin (deploy/opencode/plugin/mc2-governor.ts):
  Werkzeug-Zaun (git push, rm -rf, sudo, curl|sh — bewiesen), Pruef-Tor auf
  session.idle mit Selbstreparatur, Savepoint statt Kompression, Meldungen an
  Lucys Briefkasten mit eigenem Absender 'loop'.

Stufe 4 — Mannschaft (opencode.json):
  plan+build -> coder (51,5 t/s) · explore -> hermes (69,6 t/s, warm, gratis) ·
  review -> Devstral (15,0 t/s, FREMDE Modellfamilie gegen blinde Flecken).
  Der 63-GB-Planer faellt aus der Tagesrolle raus.

Stufe 5 — ein Regelwerk, zwei Ausloeser:
  deploy/opencode-lauf.sh faehrt dieselbe Bau-Pruef-Schleife unbeaufsichtigt
  (die Schleife liegt hier UND im Plugin: bei `opencode run` endet der Prozess,
  bevor session.idle fertig ist — gemessen). Bricht ab, wenn der Governor fehlt.
  gitea-repo-create.sh saet jetzt VERIFY neben der CI-Ampel: jedes neue Repo
  wird mit Innen- UND Aussen-Pruefung geboren.

Oberflaeche:
  Token-Waechter-Kachel im Cockpit (Fuellstand, Marken, Ehrlichkeits-Nachweis),
  /api/governor als gleichursprüngliches Fenster, Devstral im Modellkatalog,
  Rollen-Texte auf die neue Mannschaft aktualisiert.

.gitattributes: deploy/**/*.py auf LF (deploy/*.py greift nur eine Ebene tief).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-25 21:39:32 +02:00

542 lines
23 KiB
Python

#!/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 HARD_CEILING > 0 and est >= HARD_CEILING:
return "hard", body, est, streaming, model, chars
if est >= THRESHOLD and isinstance(messages, list):
if not safe_to_append(messages):
# Mitten im Werkzeug-Austausch: nicht dazwischenfunken, naechste Runde.
return "defer", body, est, streaming, model, chars
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())