""" Ideen-Queue — DIE eine Sammelstelle für Ideen des Commanders (Entscheid 10.07.2026). Die Queue selbst ist das NATIVE Hermes-Kanban (SQLite-Board, Dispatcher im Gateway): Idee rein (triage) → der Specifier arbeitet sie aus → ein Worker-Profil setzt sie um → Code-Ergebnisse kommen als Vorschlags-Branch zurück und erscheinen als Karte im Auftragsbuch. Dieser Service ist nur die MC2-Tür dazu: anlegen + anzeigen + aufräumen, alles über die Hermes-CLI (stabile --json-Schnittstelle, kein Griff in die kanban.db). Türen insgesamt: Telegram/Desktop (natives kanban_create-Tool), Lucy-Voice (mcp_voice.idee_notieren → POST /api/ideen) und die Zentrale (Auftragsbuch-Tab → hier). Lokal (Windows-Dev) ist alles harmlos: available=False, Aktionen geben Fehler statt zu crashen. """ import json import logging import os import re import sqlite3 import subprocess import time from pathlib import Path log = logging.getLogger(__name__) HERMES = Path(os.environ.get("MC_HERMES_BIN", "~/.local/bin/hermes")).expanduser() # Ketten-Sicht (20.07.): Eltern-Links + Heartbeat-Notizen stehen NICHT im `list --json` # der CLI — dafür (und NUR dafür) ein read-only-Blick in die kanban.db. Das bricht die # „alles über die CLI"-Regel bewusst minimal: zwei stabile Tabellen (task_links 2 Spalten, # task_events), Verbindung strikt mode=ro, jeder Fehler fällt lautlos auf die flache # Sicht zurück. Schreiben läuft weiterhin ausschließlich über die CLI. KANBAN_DB = Path(os.environ.get("MC_KANBAN_DB", "~/.hermes/kanban.db")).expanduser() _TASK_ID_RX = re.compile(r"^t_[0-9a-f]{4,16}$") # Die CLI kostet pro Aufruf ein paar Sekunden Python-Start — die UI pollt aber. _list_cache: dict = {"ts": 0.0, "data": None} _LIST_EVERY = 15.0 # s # Cache für worker log _log_cache: dict = {"ts": 0.0, "task_id": None, "data": None} _LOG_EVERY = 4.0 # s # Cache für Ergebnis (Abschluss-Zusammenfassung) — ändert sich nach done nie mehr _erg_cache: dict = {} _ERG_MAX = 40 # simple Deckelung gegen unbegrenztes Wachsen def _available() -> bool: return os.name == "posix" and HERMES.is_file() def _hermes(args: list[str], timeout: int = 60) -> subprocess.CompletedProcess: return subprocess.run([str(HERMES), "kanban", *args], capture_output=True, text=True, timeout=timeout) def _parse_json(stdout: str): """CLI-Ausgabe kann Log-Zeilen vor dem JSON enthalten — ab der ersten Klammer parsen.""" text = stdout or "" for opener in ("[", "{"): idx = text.find(opener) if idx >= 0: try: return json.loads(text[idx:]) except json.JSONDecodeError: continue return None def _blocker_frage(task_id: str) -> dict | None: """Für eine blockierte Aufgabe: die Frage/den Grund des Workers holen (needs_input etc.). Die Wahrheit steckt im jüngsten `blocked`-Event (payload.reason/kind); fällt das leer aus, greifen wir auf den letzten „BLOCKED:"-Kommentar zurück. Ein Extra-`show`-Aufruf pro hängender Karte ist billig — die gibt es fast nie und wenn, dann nur eine Handvoll. """ try: r = _hermes(["show", task_id, "--json"], timeout=30) d = _parse_json(r.stdout) except Exception: log.warning("ideen: kanban show %s fehlgeschlagen", task_id, exc_info=True) return None if not isinstance(d, dict): return None frage, kind = "", "" for ev in d.get("events") or []: # chronologisch → letztes blocked = aktive Frage if isinstance(ev, dict) and ev.get("kind") == "blocked": p = ev.get("payload") or {} frage = (p.get("reason") or "").strip() kind = (p.get("kind") or "").strip() if not frage: # Fallback: letzter „BLOCKED:"-Kommentar for c in d.get("comments") or []: body = (c.get("body") or "").strip() if body.startswith("BLOCKED:"): frage = body[len("BLOCKED:"):].strip() return {"frage": frage[:800], "kind": kind} def _projekt_titel(titles: list[str]) -> str: """Anzeigename einer Karten-Familie: gemeinsamer Titel-Präfix (z. B. „Homelab-Dashboard 2.0"), sonst der Titel der ältesten Karte. Reine Anzeige-Heuristik, nichts hängt davon ab.""" if not titles: return "Projekt" prefix = titles[0] for t in titles[1:]: while prefix and not t.startswith(prefix): prefix = prefix[:-1] prefix = prefix.strip(" –-:·") return prefix if len(prefix) >= 10 else titles[0][:80] def _ketten_anreichern(items: list[dict]) -> list[dict]: """Ketten-Sicht über die Queue legen (read-only, best-effort). Ergänzt pro Karte: `eltern` (offene Vorgänger-IDs), `wartet_auf` (deren Titel), `notiz` (jüngste Heartbeat-Notiz des Workers bei running) und `projekt` (Familien-Schlüssel). Liefert die Projekt-Gruppen für die UI: zusammenhängende Karten (task_links) = ein Projekt; Einzelkarten bleiben ohne `projekt`. Jeder Fehler (Windows-Dev, Schema-Drift, Lock) → unveränderte flache Sicht. """ for it in items: it.setdefault("eltern", []) it.setdefault("wartet_auf", []) it.setdefault("notiz", None) it.setdefault("projekt", None) it.setdefault("prio", 0) it.setdefault("puls_alter", None) it.setdefault("gestartet", None) if not items or not KANBAN_DB.is_file(): return [] by_id = {it["id"]: it for it in items if it.get("id")} try: con = sqlite3.connect(f"file:{KANBAN_DB}?mode=ro", uri=True, timeout=3) try: con.execute("PRAGMA busy_timeout=2000") links = [(p, c) for p, c in con.execute( "SELECT parent_id, child_id FROM task_links") if p in by_id or c in by_id] # Prio (für die Tausch-Knöpfe) + Lebenszeichen (Hänger-Warnung): ein Rutsch. ph = ",".join("?" * len(by_id)) jetzt = time.time() for tid, prio, hb, started, pid in con.execute( f"SELECT id, priority, last_heartbeat_at, started_at, worker_pid" f" FROM tasks WHERE id IN ({ph})", list(by_id)): it = by_id[tid] it["prio"] = prio or 0 if it.get("status") == "running": it["gestartet"] = started puls = hb or started it["puls_alter"] = int(jetzt - puls) if puls else None for it in items: if it.get("status") != "running": continue for (payload,) in con.execute( "SELECT payload FROM task_events WHERE task_id=? AND kind='heartbeat' " "AND payload IS NOT NULL ORDER BY created_at DESC LIMIT 5", (it["id"],)): try: note = (json.loads(payload) or {}).get("note") except Exception: note = None if note: it["notiz"] = str(note)[:300] break finally: con.close() except Exception: log.debug("ideen: Ketten-Anreicherung fehlgeschlagen", exc_info=True) return [] # Eltern/wartet_auf: nur Vorgänger zählen, die selbst noch offen in der Liste stehen — # archivierte/fremde Parents gelten als erledigt. ERLEDIGT = ("done",) for parent, child in links: it = by_id.get(child) p = by_id.get(parent) if it and p and p.get("status") not in ERLEDIGT: it["eltern"].append(parent) if len(it["wartet_auf"]) < 3: it["wartet_auf"].append((p.get("titel") or parent)[:90]) # Familien = zusammenhängende Komponenten über die Links (Union-Find, klein genug). chef: dict[str, str] = {} def boss(x: str) -> str: while chef.get(x, x) != x: chef[x] = chef.get(chef[x], chef[x]) x = chef[x] return x for parent, child in links: if parent in by_id and child in by_id: a, b = boss(parent), boss(child) if a != b: chef[a] = b familien: dict[str, list[dict]] = {} for it in items: familien.setdefault(boss(it["id"]), []).append(it) projekte = [] for wurzel, mitglieder in familien.items(): if len(mitglieder) < 2: continue mitglieder.sort(key=lambda i: i.get("erstellt") or 0) key = mitglieder[0]["id"] for it in mitglieder: it["projekt"] = key fertig = sum(1 for i in mitglieder if i["status"] in ERLEDIGT) projekte.append({ "key": key, "titel": _projekt_titel([i["titel"] for i in mitglieder]), "gesamt": len(mitglieder), "fertig": fertig, "laeuft": sum(1 for i in mitglieder if i["status"] == "running"), "haengt": sum(1 for i in mitglieder if i["status"] == "blocked"), }) # Aktive Projekte zuerst (hängend > laufend > Rest), dann nach Größe. projekte.sort(key=lambda p: (-p["haengt"], -p["laeuft"], -p["gesamt"])) return projekte def list_queue() -> dict: """Alle nicht archivierten Aufgaben des Boards, jüngste zuerst — die Queue-Sicht der UI.""" if not _available(): return {"available": False, "items": [], "offen": 0} now = time.time() if _list_cache["data"] is not None and now - _list_cache["ts"] < _LIST_EVERY: return _list_cache["data"] try: r = _hermes(["list", "--json", "--sort", "created-desc"]) raw = _parse_json(r.stdout) except Exception: log.warning("ideen: kanban list fehlgeschlagen", exc_info=True) raw = None if not isinstance(raw, list): # CLI kaputt/Timeout → ehrlich leer melden, aber nicht cachen (nächster Poll versucht's neu) return {"available": True, "items": [], "offen": 0, "fehler": "Queue nicht lesbar"} items = [] for t in raw[:60]: try: items.append({ "id": t.get("id", ""), "titel": (t.get("title") or "(ohne Titel)")[:200], "body": (t.get("body") or "")[:600], "status": t.get("status", "?"), "assignee": t.get("assignee"), # IDE-Welt = das Vorbereiter-Profil (Konzept + Repo + Doku, kein Bau). "welt": "ide" if t.get("assignee") == "projektstart" else "box", "erstellt": t.get("created_at"), "fertig": t.get("completed_at"), "von": t.get("created_by"), "frage": None, # bei blocked die Worker-Frage (unten nachgeladen) "frage_kind": None, }) except Exception: continue # Hängende Karten anreichern: die Frage des Workers holen, damit die Zentrale sie # beantworten kann (Sackgasse-Fix). Deckel gegen Ausreißer, damit ein Poll nie hängt. for it in [i for i in items if i["status"] == "blocked"][:6]: info = _blocker_frage(it["id"]) if info: it["frage"] = info["frage"] it["frage_kind"] = info["kind"] projekte = _ketten_anreichern(items) offen = sum(1 for i in items if i["status"] not in ("done",)) data = {"available": True, "items": items, "offen": offen, "projekte": projekte} _list_cache.update(ts=now, data=data) return data def add_idea(titel: str, notiz: str = "", created_by: str = "mc2-ui", welt: str = "box") -> dict: """Idee in die Queue legen. welt='box' (Default): triage — der Specifier arbeitet sie aus, die Werkstatt baut, Ergebnis ist ein Patch-Vorschlag im Auftragsbuch (Box-Welt). welt='ide': geht DIREKT an das Profil `projektstart` (am Triage/Auto-Zerleger vorbei) — die Box härtet ein Konzept (konzept-fliessband), legt ein leeres Gitea-Repo an und füllt es NUR mit Doku. Ergebnis ist ein vorbereitetes Projekt, das der Commander in Zed selbst baut (IDE-Welt). Bewusst KEIN --triage: so fasst der Auto-Zerleger die Karte nie an und die Werkstatt baut sie nie (der Fan-out-Fehler vom 20.07.).""" if not _available(): return {"ok": False, "error": "Die Ideen-Queue lebt auf der Box."} titel = (titel or "").strip() if not (3 <= len(titel) <= 200): return {"ok": False, "error": "Titel bitte zwischen 3 und 200 Zeichen."} welt = "ide" if str(welt).lower() == "ide" else "box" if welt == "ide": args = ["create", titel, "--assignee", "projektstart", "--created-by", created_by, "--json"] else: args = ["create", titel, "--triage", "--created-by", created_by, "--json"] notiz = (notiz or "").strip() if notiz: args[2:2] = ["--body", notiz[:4000]] try: r = _hermes(args) except Exception as exc: return {"ok": False, "error": f"Queue nicht erreichbar: {exc}"} task = _parse_json(r.stdout) if r.returncode != 0 or not isinstance(task, dict) or not task.get("id"): return {"ok": False, "error": (r.stderr or r.stdout or "kanban create fehlgeschlagen").strip()[:300]} _list_cache["ts"] = 0.0 # nächster Poll zeigt die neue Idee sofort try: from services import announce if welt == "ide": announce.add(f"Neues IDE-Projekt in Vorbereitung: „{titel}“ ({task['id']}) — die Box " "härtet das Konzept und legt das Repo an, dann baust du in Zed.", "[Ideen-Queue]", "ideen-queue", "silent") else: announce.add(f"Neue Idee in der Queue: „{titel}“ ({task['id']}) — die Box arbeitet sie aus.", "[Ideen-Queue]", "ideen-queue", "silent") except Exception: pass return {"ok": True, "id": task["id"], "status": task.get("status"), "welt": welt} def answer_task(task_id: str, antwort: str) -> dict: """Eine hängende Aufgabe beantworten: Antwort als Kommentar protokollieren + entsperren. `unblock --reason` schreibt die Antwort als „UNBLOCK:"-Kommentar und stellt die Karte auf ready — der Dispatcher spawnt den Worker neu, der die Antwort im Task-Kontext vorfindet und weitermacht. Genau der Ausweg aus der Sackgasse „Aufgabe fragt, aber niemand kann antworten". """ if not _available(): return {"ok": False, "error": "Die Ideen-Queue lebt auf der Box."} if not _TASK_ID_RX.match(task_id or ""): return {"ok": False, "error": f"Keine gültige Aufgaben-Nummer: {task_id!r}"} antwort = (antwort or "").strip() if not (1 <= len(antwort) <= 2000): return {"ok": False, "error": "Bitte eine Antwort zwischen 1 und 2000 Zeichen."} try: r = _hermes(["unblock", task_id, "--reason", antwort]) except Exception as exc: return {"ok": False, "error": f"Queue nicht erreichbar: {exc}"} if r.returncode != 0: return {"ok": False, "error": (r.stderr or r.stdout or "Entsperren fehlgeschlagen").strip()[:300]} _list_cache["ts"] = 0.0 # nächster Poll zeigt den neuen Status (ready/running) sofort try: from services import announce announce.add(f"Deine Antwort ging an die hängende Aufgabe {task_id} — die Box macht weiter.", "[Ideen-Queue]", "ideen-queue", "silent") except Exception: pass return {"ok": True} def archive_task(task_id: str) -> dict: """Erledigtes/Verworfenes aus der Sicht räumen (Kanban-Archiv, nichts wird gelöscht).""" if not _available(): return {"ok": False, "error": "Die Ideen-Queue lebt auf der Box."} if not _TASK_ID_RX.match(task_id or ""): return {"ok": False, "error": f"Keine gültige Aufgaben-Nummer: {task_id!r}"} try: r = _hermes(["archive", task_id]) except Exception as exc: return {"ok": False, "error": f"Queue nicht erreichbar: {exc}"} if r.returncode != 0: return {"ok": False, "error": (r.stderr or r.stdout or "Archivieren fehlgeschlagen").strip()[:300]} _list_cache["ts"] = 0.0 return {"ok": True} def ergebnis_of(task_id: str) -> dict: """Das ERGEBNIS einer fertigen Aufgabe — die Abschluss-Zusammenfassung des Workers. User-Wunsch 15.07.: „im Auftragsbuch werden Dinge getestet, ich sehe aber das Ergebnis nicht." Die Wahrheit liegt beim kanban_complete-Aufruf (summary/result) im jüngsten `completed`-Event; `list --json` trägt sie nicht → lazy per `show`, nur wenn der User die Karte aufklappt. Fertige Ergebnisse ändern sich nie → Cache. """ if not _available(): return {"available": False, "summary": "", "result": ""} if not _TASK_ID_RX.match(task_id or ""): return {"available": True, "summary": "", "result": ""} if task_id in _erg_cache: return _erg_cache[task_id] try: r = _hermes(["show", task_id, "--json"], timeout=30) d = _parse_json(r.stdout) except Exception: log.warning("ideen: kanban show %s (ergebnis) fehlgeschlagen", task_id, exc_info=True) return {"available": True, "summary": "", "result": ""} summary, result = "", "" if isinstance(d, dict): if isinstance(d.get("result"), str): result = d["result"].strip() for ev in d.get("events") or []: # chronologisch → letztes completed gewinnt if isinstance(ev, dict) and ev.get("kind") == "completed": s = ((ev.get("payload") or {}).get("summary") or "").strip() if s: summary = s data = {"available": True, "summary": summary[:4000], "result": result[:4000]} if len(_erg_cache) >= _ERG_MAX: _erg_cache.pop(next(iter(_erg_cache))) _erg_cache[task_id] = data return data # ── Karten-Steuerung (Stopp / Neuer Versuch / Prio) — User-Wunsch 20.07. ───── # Die CLI kann blocken/entsperren, aber weder haengende Worker beenden noch die # Prioritaet aendern → fuer genau diese zwei Faelle eine eng begrenzte Direkt- # Schreibstelle in die kanban.db (BEGIN IMMEDIATE, nur tasks-Spalten). Alles # andere laeuft weiter ueber die CLI. def _task_row(task_id: str) -> dict | None: try: con = sqlite3.connect(f"file:{KANBAN_DB}?mode=ro", uri=True, timeout=3) try: con.row_factory = sqlite3.Row r = con.execute("SELECT id, status, worker_pid, priority FROM tasks WHERE id=?", (task_id,)).fetchone() return dict(r) if r else None finally: con.close() except Exception: return None def _kill_tree(pid: int) -> None: """Worker samt Kind-Prozessen beenden (TERM, dann KILL). Ohne das bleiben Test-Server (uvicorn & Co.) als Waisen zurueck — live gesehen 20.07.""" import signal def kinder(p: int) -> list[int]: try: out = subprocess.run(["pgrep", "-P", str(p)], capture_output=True, text=True, timeout=5) return [int(x) for x in out.stdout.split()] except Exception: return [] def baum(p: int) -> list[int]: alle = [p] for k in kinder(p): alle.extend(baum(k)) return alle pids = baum(pid) for p in reversed(pids): try: os.kill(p, signal.SIGTERM) except OSError: pass time.sleep(2) for p in reversed(pids): try: os.kill(p, signal.SIGKILL) except OSError: pass def _requeue(task_id: str) -> None: """Karte sauber zurueck auf todo (Claim/PID/Zaehler leeren) — der Dispatcher spawnt sie beim naechsten Tick frisch (Caps gelten).""" con = sqlite3.connect(str(KANBAN_DB), timeout=15) try: con.execute("PRAGMA busy_timeout=10000") con.execute("BEGIN IMMEDIATE") con.execute( "UPDATE tasks SET status='todo', claim_lock=NULL, claim_expires=NULL," " worker_pid=NULL, current_run_id=NULL, last_heartbeat_at=NULL," " started_at=NULL, consecutive_failures=0 WHERE id=?", (task_id,)) con.commit() finally: con.close() def stop_task(task_id: str) -> dict: """Karte gezielt anhalten: generischer Block (bleibt liegen, kein Auto-Unblock- Kind) + laufenden Worker-Prozessbaum beenden.""" if not _available(): return {"ok": False, "error": "Die Ideen-Queue lebt auf der Box."} if not _TASK_ID_RX.match(task_id or ""): return {"ok": False, "error": f"Keine gültige Aufgaben-Nummer: {task_id!r}"} row = _task_row(task_id) if not row: return {"ok": False, "error": "Karte nicht gefunden."} try: r = _hermes(["block", task_id, "Vom", "Commander", "gestoppt"]) except Exception as exc: return {"ok": False, "error": f"Queue nicht erreichbar: {exc}"} if r.returncode != 0: return {"ok": False, "error": (r.stderr or r.stdout or "Blocken fehlgeschlagen").strip()[:300]} if row.get("status") == "running" and row.get("worker_pid"): _kill_tree(int(row["worker_pid"])) _list_cache["ts"] = 0.0 return {"ok": True} def retry_task(task_id: str) -> dict: """Neuer Versuch: haengende/laufende Karte → Worker-Baum beenden + frisch auf todo; blockierte Karte → entsperren. Der Dispatcher uebernimmt den Rest (Caps gelten).""" if not _available(): return {"ok": False, "error": "Die Ideen-Queue lebt auf der Box."} if not _TASK_ID_RX.match(task_id or ""): return {"ok": False, "error": f"Keine gültige Aufgaben-Nummer: {task_id!r}"} row = _task_row(task_id) if not row: return {"ok": False, "error": "Karte nicht gefunden."} status = row.get("status") try: if status == "blocked": r = _hermes(["unblock", task_id, "--reason", "Neuer Versuch vom Commander"]) if r.returncode != 0: return {"ok": False, "error": (r.stderr or r.stdout or "Entsperren fehlgeschlagen").strip()[:300]} else: if status == "running" and row.get("worker_pid"): _kill_tree(int(row["worker_pid"])) _requeue(task_id) try: _hermes(["comment", task_id, "Neuer", "Versuch", "vom", "Commander", "(Karte zurueckgesetzt)"], timeout=20) except Exception: pass except Exception as exc: return {"ok": False, "error": f"Neuer Versuch fehlgeschlagen: {exc}"} _list_cache["ts"] = 0.0 try: from services import announce announce.add(f"Karte {task_id} wurde auf Commander-Wunsch neu angestoßen.", "[Ideen-Queue]", "ideen-queue", "silent") except Exception: pass return {"ok": True} def prio_task(task_id: str, richtung: str) -> dict: """Prio hoch/runter (Dispatcher zieht ready-Karten nach `priority DESC, created ASC` → hoch = frueher dran). Direkt-Schreibstelle, die CLI kennt keine Prio-Aenderung.""" if not _available(): return {"ok": False, "error": "Die Ideen-Queue lebt auf der Box."} if not _TASK_ID_RX.match(task_id or ""): return {"ok": False, "error": f"Keine gültige Aufgaben-Nummer: {task_id!r}"} if richtung not in ("hoch", "runter"): return {"ok": False, "error": "richtung muss 'hoch' oder 'runter' sein."} delta = 1 if richtung == "hoch" else -1 try: con = sqlite3.connect(str(KANBAN_DB), timeout=15) try: con.execute("PRAGMA busy_timeout=10000") con.execute("BEGIN IMMEDIATE") cur = con.execute("UPDATE tasks SET priority = COALESCE(priority,0) + ? WHERE id=?", (delta, task_id)) con.commit() if cur.rowcount == 0: return {"ok": False, "error": "Karte nicht gefunden."} finally: con.close() except Exception as exc: return {"ok": False, "error": f"Prio-Änderung fehlgeschlagen: {exc}"} _list_cache["ts"] = 0.0 return {"ok": True} def log_of(task_id: str) -> dict: """Worker-Log einer Idee (hermes kanban log ), ANSI entfernt, letzte ~120 Zeilen.""" if not _available(): return {"available": False, "lines": []} if not _TASK_ID_RX.match(task_id or ""): return {"available": True, "lines": []} now = time.time() if ( _log_cache["data"] is not None and _log_cache["task_id"] == task_id and now - _log_cache["ts"] < _LOG_EVERY ): return _log_cache["data"] try: r = _hermes(["log", task_id]) except Exception as exc: log.warning("ideen: kanban log fehlgeschlagen", exc_info=True) return {"available": True, "lines": []} if r.returncode != 0: log.warning("ideen: kanban log exit=%s stderr=%s", r.returncode, r.stderr) return {"available": True, "lines": []} raw = r.stdout or "" # ANSI-Steuerzeichen entfernen ansi = re.compile(r"\x1b\[[0-9;]*m") cleaned = ansi.sub("", raw) # Letzte ~120 Zeilen lines = cleaned.strip().splitlines()[-120:] data = {"available": True, "lines": lines} _log_cache.update(ts=now, task_id=task_id, data=data) return data