""" 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)