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>
This commit is contained in:
+213
-92
@@ -1,34 +1,52 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Governor — duenner, zustandsloser Token-Waechter-Proxy (Phase 0).
|
||||
"""Governor v2 — Token-Waechter-Proxy vor dem MC2-Gateway.
|
||||
|
||||
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:
|
||||
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 >= 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.
|
||||
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 zustandslos: jede Anfrage wird fuer sich bewertet; keine Sitzungs-DB.
|
||||
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 Modell-Endpoint (Default http://127.0.0.1:8080)
|
||||
GOV_THRESHOLD Soft-Schwelle fuer den Einschub (Default 25000)
|
||||
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 Heuristik Zeichen->Token (Default 3.5, kalibriert)
|
||||
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 (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)
|
||||
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
|
||||
@@ -37,6 +55,7 @@ import os
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
from collections import deque
|
||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||
from urllib.parse import urlparse
|
||||
|
||||
@@ -44,16 +63,22 @@ from urllib.parse import urlparse
|
||||
|
||||
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.
|
||||
# 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.5"))
|
||||
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 "
|
||||
@@ -76,12 +101,8 @@ DEFAULT_HARDSTOP = (
|
||||
)
|
||||
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
|
||||
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."
|
||||
@@ -100,9 +121,21 @@ HOP_BY_HOP = {
|
||||
_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)
|
||||
_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)."""
|
||||
@@ -117,28 +150,77 @@ def log(line: str) -> None:
|
||||
pass
|
||||
|
||||
|
||||
def estimate_tokens(messages) -> int:
|
||||
"""Grobe, aber stabile Heuristik: Zeichen aller Nachrichteninhalte / CPT.
|
||||
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)
|
||||
|
||||
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.
|
||||
|
||||
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.
|
||||
"""
|
||||
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)
|
||||
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. Laeuft im Hintergrund-Thread."""
|
||||
"""POSTet die Meldung an die MC2-Announce-Pipeline (Lucy spricht sie)."""
|
||||
try:
|
||||
text = (ANNOUNCE_TEXT.replace("{est}", str(est))
|
||||
.replace("{threshold}", str(THRESHOLD)))
|
||||
@@ -158,8 +240,7 @@ def _post_announce(est) -> None:
|
||||
|
||||
|
||||
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."""
|
||||
"""Sprach-Signal an Lucy — gedrosselt (eine Aeusserung je Episode)."""
|
||||
if not _an:
|
||||
return
|
||||
global _last_announce
|
||||
@@ -171,14 +252,35 @@ def maybe_announce(est) -> None:
|
||||
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/0.2"
|
||||
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):
|
||||
@@ -194,6 +296,19 @@ class Handler(BaseHTTPRequestHandler):
|
||||
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:
|
||||
@@ -206,8 +321,6 @@ class Handler(BaseHTTPRequestHandler):
|
||||
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")
|
||||
|
||||
@@ -215,19 +328,23 @@ class Handler(BaseHTTPRequestHandler):
|
||||
est = None
|
||||
streaming = False
|
||||
model = ""
|
||||
chars = 0
|
||||
if is_chat and body:
|
||||
action, body, est, streaming, model = self._decide(body)
|
||||
action, body, est, streaming, model, chars = self._decide(body)
|
||||
if action in ("soft", "hard"):
|
||||
maybe_announce(est) # Sprach-Signal an Lucy (gedrosselt)
|
||||
maybe_announce(est)
|
||||
|
||||
# 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")
|
||||
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
|
||||
|
||||
# Header fuer Upstream aufbereiten.
|
||||
out_headers = {}
|
||||
for k, v in self.headers.items():
|
||||
kl = k.lower()
|
||||
@@ -242,7 +359,7 @@ class Handler(BaseHTTPRequestHandler):
|
||||
|
||||
conn = None
|
||||
try:
|
||||
conn = http.client.HTTPConnection(UP_HOST, UP_PORT, timeout=600)
|
||||
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:
|
||||
@@ -255,7 +372,6 @@ class Handler(BaseHTTPRequestHandler):
|
||||
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)
|
||||
@@ -266,8 +382,7 @@ class Handler(BaseHTTPRequestHandler):
|
||||
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.
|
||||
# Durchreichen); read() wuerde puffern und Streaming haengen lassen.
|
||||
chunk = resp.read1(65536)
|
||||
if not chunk:
|
||||
break
|
||||
@@ -283,41 +398,55 @@ class Handler(BaseHTTPRequestHandler):
|
||||
|
||||
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 "ok"
|
||||
log(f"chat est={est} exact={exact_s} thr={THRESHOLD} {flag} "
|
||||
f"stream={streaming} status={resp.status}")
|
||||
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: passthrough | soft (Einschub) | hard (Selbstantwort).
|
||||
|
||||
Rueckgabe: (action, body, est, streaming, model).
|
||||
"""
|
||||
"""Aktion bestimmen. Rueckgabe: (action, body, est, streaming, model, chars)."""
|
||||
try:
|
||||
data = json.loads(body)
|
||||
except (ValueError, UnicodeDecodeError):
|
||||
return "passthrough", body, None, False, ""
|
||||
return "passthrough", body, None, False, "", 0
|
||||
if not isinstance(data, dict):
|
||||
return "passthrough", body, None, False, ""
|
||||
return "passthrough", body, None, False, "", 0
|
||||
|
||||
messages = data.get("messages")
|
||||
streaming = bool(data.get("stream"))
|
||||
model = data.get("model", "") or ""
|
||||
est = estimate_tokens(messages if isinstance(messages, list) else [])
|
||||
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
|
||||
return "hard", body, est, streaming, model, chars
|
||||
|
||||
if est >= THRESHOLD and isinstance(messages, list):
|
||||
# Sichere Substitution statt str.format: ein operator-gesetzter
|
||||
# GOV_DIRECTIVE mit { } (JSON/Code-Beispiel) darf nicht crashen.
|
||||
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
|
||||
return "soft", json.dumps(data).encode("utf-8"), est, streaming, model, chars
|
||||
|
||||
return "passthrough", body, est, streaming, model
|
||||
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)."""
|
||||
@@ -366,16 +495,7 @@ class Handler(BaseHTTPRequestHandler):
|
||||
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
|
||||
self._send_json(status, {"error": msg})
|
||||
|
||||
@staticmethod
|
||||
def _scan_prompt_tokens(tail: bytearray):
|
||||
@@ -405,8 +525,9 @@ def main() -> int:
|
||||
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}")
|
||||
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:
|
||||
|
||||
Reference in New Issue
Block a user