Files
mission-control-v2/backend/services/ketten_digest.py
T
Hitonabi 8153d14d9f Queue-Steuerung: Stopp/Neuer-Versuch/Prio + Haenger-Warnung (UI+Telegram)
Karten lassen sich gezielt anhalten (generischer Block + Prozessbaum-Kill),
frisch starten (Retry beendet Worker samt Kindern - uvicorn-Falle) und per
Prio-Pfeilen umsortieren (Dispatcher zieht priority DESC nativ). Laufende
Karten ohne Heartbeat >10 min bekommen ein Warn-Badge; der Ketten-Digest
alarmiert nach 15 min per Telegram, inkl. Kontext-Verdichtungs-Hinweis.
Ausloeser: Scaffold-Worker hing 45 min an selbst gestartetem uvicorn.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-20 13:26:09 +02:00

192 lines
7.9 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
Ketten-Digest — Telegram-Status für lange Aufgaben-Ketten (User-Wunsch 20.07.2026).
Wenn die Werkstatt eine Karten-Familie abarbeitet (z. B. „Homelab-Dashboard 2.0" mit
19 verketteten Karten), soll der Commander nicht raten müssen: dieser Wächter liest
die Queue-Sicht (services/ideen.py, inkl. Ketten-Anreicherung) und meldet per Telegram
• SOFORT, wenn eine Karte hängt und eine Frage an den Commander hat (blocked) —
einmal pro Karte, mit der Frage im Wortlaut,
• bei MEILENSTEINEN (Projekt komplett fertig) sofort,
• sonst höchstens einmal pro PULS_SEKUNDEN (Default 60 min) einen Zwischenstand
je aktivem Projekt („5/19 fertig · läuft: … · als Nächstes: …").
Taktung war explizite User-Wahl („Meilensteine + 60-min-Puls, Fragen immer sofort").
Zustand (was wurde wann gemeldet) liegt als JSON neben den Modellen und übersteht
Neustarts — sonst käme nach jedem Deploy ein Duplikat-Schwall. Auf Windows (Dev) No-op.
"""
import asyncio
import json
import logging
import os
import time
import zlib
from pathlib import Path
from config import MODELS_DIR
log = logging.getLogger(__name__)
ENABLED = os.name == "posix" and os.environ.get("MC_KETTEN_DIGEST", "1") != "0"
INTERVAL = int(os.environ.get("MC_KETTEN_DIGEST_INTERVAL", "300")) # Prüf-Tick: 5 min
PULS_SEKUNDEN = int(os.environ.get("MC_KETTEN_DIGEST_PULS", "3600")) # Zwischenstand: 60 min
# Hänger-Alarm: laufende Karte ohne Lebenszeichen (Heartbeat) länger als diese Schwelle
# → sofortige Telegram-Warnung (einmal pro Karten-Lauf). Live-Fall 20.07.: Worker wartete
# 45 min auf einen selbst gestarteten uvicorn — „pid_alive" hielt den Claim ewig frisch.
HANG_SEKUNDEN = int(os.environ.get("MC_KETTEN_HANG", "900"))
KANBAN_LOGS = Path(os.environ.get("MC_KANBAN_LOGS", "~/.hermes/kanban/logs")).expanduser()
STATE_PATH = Path(os.environ.get("MC_KETTEN_DIGEST_STATE", str(MODELS_DIR / "mc2-ketten-digest.json")))
def _kontext_druck(task_id: str) -> int:
"""Wie oft der Worker zuletzt den Kontext verdichten musste (Log-Marker) —
mehrfaches „Compacting context" heißt: die Aufgabe sprengt das Fenster."""
try:
p = KANBAN_LOGS / f"{task_id}.log"
with p.open("rb") as f:
f.seek(max(0, p.stat().st_size - 200_000))
tail = f.read().decode("utf-8", errors="replace")
return tail.count("Compacting context")
except Exception:
return 0
def _load_state() -> dict:
try:
d = json.loads(STATE_PATH.read_text(encoding="utf-8"))
return d if isinstance(d, dict) else {}
except Exception:
return {}
def _save_state(state: dict) -> None:
try:
tmp = STATE_PATH.with_suffix(".tmp")
tmp.write_text(json.dumps(state, ensure_ascii=False), encoding="utf-8")
tmp.replace(STATE_PATH)
except OSError:
log.warning("ketten-digest: Zustand %s nicht schreibbar", STATE_PATH, exc_info=True)
def _melden(subject: str, text: str) -> None:
"""Aufs Handy UND in den Briefkasten (silent — Lucy muss den Status nicht sprechen)."""
from services import announce
announce.notify_telegram(subject, text)
try:
announce.add(text, subject, "ketten-digest", "silent")
except Exception:
pass
def _tick() -> None:
from services import ideen
data = ideen.list_queue()
if not data.get("available"):
return
items = data.get("items") or []
projekte = data.get("projekte") or []
by_id = {i["id"]: i for i in items if i.get("id")}
state = _load_state()
gemeldete_fragen: dict = state.setdefault("fragen", {})
projekt_state: dict = state.setdefault("projekte", {})
dirty = False
now = int(time.time())
# 1) Hängende Karten mit Frage → sofort, einmal pro Karte. Der Schlüssel enthält den
# Fragen-Text-Hash: hängt dieselbe Karte später mit NEUER Frage, wird wieder gemeldet.
for it in items:
if it.get("status") != "blocked":
continue
frage = (it.get("frage") or "").strip()
key = f"{it['id']}:{zlib.crc32(frage.encode('utf-8')):x}"
if key in gemeldete_fragen:
continue
text = (f"Aufgabe hängt und wartet auf dich: „{(it.get('titel') or '')[:120]}\n"
+ (f"Frage: {frage[:500]}\n" if frage else "")
+ "Antworten geht im Auftragsbuch (Zentrale) — die Box macht dann weiter.")
_melden("[Werkstatt]", text)
gemeldete_fragen[key] = now
dirty = True
# 1b) Hänger-Alarm: laufende Karte ohne Lebenszeichen → einmal pro Karten-Lauf warnen.
gemeldete_haenger: dict = state.setdefault("haenger", {})
for it in items:
if it.get("status") != "running":
continue
puls = it.get("puls_alter")
if puls is None or puls < HANG_SEKUNDEN:
continue
key = f"{it['id']}:{it.get('gestartet') or 0}"
if key in gemeldete_haenger:
continue
minuten = int(puls / 60)
text = (f"Karte scheint zu hängen: „{(it.get('titel') or '')[:120]}“ — "
f"seit {minuten} min kein Lebenszeichen vom Worker.")
druck = _kontext_druck(it["id"])
if druck >= 2:
text += (f"\nDer Worker musste {druck}× den Kontext verdichten — "
"die Aufgabe ist womöglich zu groß geschnitten.")
text += "\nIn der Zentrale: „Neuer Versuch“ startet sie frisch, „Stopp“ hält sie an."
_melden("[Werkstatt]", text)
gemeldete_haenger[key] = now
dirty = True
if len(gemeldete_haenger) > 200:
for k in sorted(gemeldete_haenger, key=gemeldete_haenger.get)[:100]:
gemeldete_haenger.pop(k, None)
dirty = True
# Fragen-Gedächtnis klein halten (Karten verschwinden irgendwann ins Archiv).
if len(gemeldete_fragen) > 200:
for k in sorted(gemeldete_fragen, key=gemeldete_fragen.get)[:100]:
gemeldete_fragen.pop(k, None)
dirty = True
# 2) Projekt-Status: Abschluss sofort, sonst gedrosselter Puls solange gearbeitet wird.
for p in projekte:
ps = projekt_state.setdefault(p["key"], {})
fertig, gesamt = p.get("fertig", 0), p.get("gesamt", 0)
if gesamt > 0 and fertig >= gesamt and not ps.get("abschluss_gemeldet"):
_melden("[Werkstatt]", f"Projekt fertig: „{p['titel']}“ — alle {gesamt} Karten erledigt. "
"Ergebnisse liegen im Auftragsbuch.")
ps.update(abschluss_gemeldet=True, letzter_puls=now, letzter_stand=fertig)
dirty = True
continue
aktiv = p.get("laeuft", 0) > 0
if not aktiv:
continue
if now - int(ps.get("letzter_puls") or 0) < PULS_SEKUNDEN:
continue
laufende = [i for i in items if i.get("projekt") == p["key"] and i["status"] == "running"]
naechste = [i for i in items if i.get("projekt") == p["key"]
and i["status"] in ("ready", "todo") and not i.get("wartet_auf")]
zeilen = [f"Werkstatt-Status „{p['titel']}“: {fertig}/{gesamt} fertig."]
for l in laufende[:2]:
note = f"{l['notiz']}" if l.get("notiz") else ""
zeilen.append(f"Läuft: {l['titel'][:90]}{note}")
if naechste:
zeilen.append(f"Als Nächstes: {naechste[0]['titel'][:90]}")
haengt = p.get("haengt", 0)
if haengt:
zeilen.append(f"{haengt} Karte(n) warten auf deine Antwort im Auftragsbuch.")
_melden("[Werkstatt]", "\n".join(zeilen))
ps.update(letzter_puls=now, letzter_stand=fertig)
dirty = True
if dirty:
_save_state(state)
async def digest_loop() -> None:
"""Hintergrund-Task im MC2-Lifespan (Muster: reminders_loop)."""
log.info("ketten-digest: aktiv (Tick %ss, Puls %ss)", INTERVAL, PULS_SEKUNDEN)
while True:
try:
await asyncio.to_thread(_tick)
except Exception:
log.debug("ketten-digest: Tick fehlgeschlagen", exc_info=True)
await asyncio.sleep(INTERVAL)