"""Token-Erfassung für den Builtin-Gateway. Parst die `usage`-Felder aus llama-swap-Antworten (Stream + Non-Stream) und meldet sie an token_stats. Hält den gateway_proxy-Router dünn und ersetzt die zuvor inline verstreute, still scheiternde String-Suche durch einen testbaren SSE-Zeilenparser. """ import json import logging import time from services.token_stats import increment_tokens log = logging.getLogger(__name__) # Kontext-Limit-Warnung (User-Wunsch 15.07.: „es sollte eine Warnung geben — Kontext ist # nicht unendlich"): finish_reason=length heißt, eine Antwort ist am Token-/Kontext-Budget # ABGERISSEN — bei Nacht-Workern stirbt damit still die halbe Arbeit (4 Worker am 15.07.). # Statt still: Eintrag in den Briefkasten (silent → Chronik/Panel, kein Sprach-Spam), # je Modell höchstens alle 10 min. Läuft im mc2-gateway-Prozess → announce liefert per # MC_ANNOUNCE_HTTP beim Steuerpult ab (Unit-Env), nie direkt in die Store-Datei. _trunc_last: dict[str, float] = {} _TRUNC_EVERY = 600.0 # s def warn_truncation(model: str) -> None: now = time.time() if now - _trunc_last.get(model, 0.0) < _TRUNC_EVERY: return _trunc_last[model] = now log.warning("gateway: finish_reason=length bei %s — Antwort am Kontext-/Token-Limit abgerissen", model) try: from services import announce announce.add( f"Eine Antwort von „{model}“ ist am Token-/Kontext-Limit abgerissen " f"(finish_reason=length). War das ein Nacht-Worker, ist seine Aufgabe zu groß " f"geschnitten — besser in Etappen teilen (eine Etappe = ein Worker-Lauf).", "[Kontext-Limit]", "gateway", "silent") except Exception: log.warning("gateway: Kontext-Limit-Warnung nicht zustellbar", exc_info=True) def record_usage(usage: dict | None, model: str) -> None: """Ein usage-Objekt verbuchen (no-op bei None/leer).""" if not usage: return prompt = usage.get("prompt_tokens", 0) completion = usage.get("completion_tokens", 0) if prompt or completion: increment_tokens(prompt, completion, model=model) def record_stream_chunk(chunk: bytes, model: str) -> None: """Rohen SSE-Chunk auf `usage` prüfen und Tokens verbuchen. Fehler werden geloggt (debug) statt verschluckt — ein defekter Chunk bricht den Stream nicht.""" if b'"finish_reason":"length"' in chunk or b'"finish_reason": "length"' in chunk: warn_truncation(model) if b'"usage"' not in chunk: return text = chunk.decode("utf-8", errors="ignore") for line in text.splitlines(): if not line.startswith("data:"): continue data_str = line[5:].strip() if not data_str or data_str == "[DONE]": continue try: record_usage(json.loads(data_str).get("usage"), model) except json.JSONDecodeError: log.debug("gateway stream: usage-Parsing fehlgeschlagen: %s", data_str[:120])