7b9752261a
Der Worker schreibt die Repo-URL in der Abschluss-Zusammenfassung als
`https://...` — die Lookahead-Liste der Regex kannte den Backtick nicht,
also fand die Konzept-Ansicht kein Repo ('wird gerade vorbereitet') obwohl
Repo + KONZEPT.md laengst da waren. Fix: greedy matchen, Treffer hinterher
saeubern (Backticks/Satzzeichen/.git) statt das Endezeichen zu raten.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
724 lines
31 KiB
Python
724 lines
31 KiB
Python
"""
|
||
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
|
||
|
||
|
||
# Aufwand-Text je Tiefe, den der Commander per Knopf waehlt — wird dem Auftrag
|
||
# vorangestellt und ueberstimmt die Groessen-Selbsteinschaetzung von konzept-fliessband.
|
||
_AUFWAND = {
|
||
"simpel": ("AUFWAND (vom Commander gewaehlt): SIMPEL — nutze die KLEINE konzept-fliessband-"
|
||
"Variante: 2 Fach-Rollen, 1 Denk-Runde + 1 Haertetest, frueh raus bei WASSERDICHT. "
|
||
"Das ueberstimmt deine eigene Groessen-Einschaetzung."),
|
||
"gruendlich": ("AUFWAND (vom Commander gewaehlt): GRUENDLICH — nutze die VOLLE konzept-fliessband-"
|
||
"Kaskade: 3-5 Fach-Rollen, bis 3 Denk-Runden, mehrfacher Advocatus-Diaboli-Haertetest. "
|
||
"Das ueberstimmt deine eigene Groessen-Einschaetzung, auch wenn das Projekt klein wirkt."),
|
||
}
|
||
|
||
|
||
def add_idea(titel: str, notiz: str = "", created_by: str = "mc2-ui",
|
||
welt: str = "box", tiefe: str = "") -> 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.).
|
||
tiefe (nur ide): 'simpel' | 'gruendlich' — vom Knopf gewaehlte Konzept-Tiefe; leer =
|
||
die Box schaetzt selbst (Fallback fuer Telegram/Lucy)."""
|
||
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"
|
||
notiz = (notiz or "").strip()
|
||
if welt == "ide":
|
||
args = ["create", titel, "--assignee", "projektstart", "--created-by", created_by, "--json"]
|
||
aufwand = _AUFWAND.get(str(tiefe).lower())
|
||
if aufwand:
|
||
notiz = aufwand + ("\n\n" + notiz if notiz else "")
|
||
else:
|
||
args = ["create", titel, "--triage", "--created-by", created_by, "--json"]
|
||
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}
|
||
|
||
|
||
# ── Konzept-Ansicht (IDE-Welt) — KONZEPT.md aus dem vorbereiteten Repo holen ──
|
||
# Der projektstart-Worker legt ein Gitea-Repo an und schreibt das wasserdichte
|
||
# Konzept als KONZEPT.md hinein. Hier holen wir es (Gitea-API + Auth aus
|
||
# ~/.git-credentials), damit der Commander es DIREKT im Auftragsbuch liest,
|
||
# ohne das Repo aufzumachen (User-Wunsch 20.07.).
|
||
_GITEA_HOST = "git.tobisniceshomelab.ddnsfree.com"
|
||
# Kein Raten am Zeilenende: greedy bis zum ersten Nicht-Repo-Zeichen matchen und den
|
||
# Treffer HINTERHER säubern. (Der Worker schreibt die URL gern in Markdown-Backticks —
|
||
# eine Lookahead-Liste erwischt sowas nie zuverlässig, Bug-Fund 20.07.)
|
||
_GITEA_REPO_RX = re.compile(
|
||
re.escape(_GITEA_HOST) + r"/([A-Za-z0-9_.-]+/[A-Za-z0-9_][A-Za-z0-9_.-]*)"
|
||
)
|
||
_REPO_FALLBACK_RX = re.compile(r"[Rr]epo\S*[:\s]+[`'\"]?([A-Za-z0-9_.-]+/[A-Za-z0-9_][A-Za-z0-9_.-]*)")
|
||
|
||
|
||
def _clean_repo(raw: str) -> str | None:
|
||
"""Trailing Markdown-/Satzzeichen und .git vom Treffer abschneiden."""
|
||
r = (raw or "").strip().strip("`'\"")
|
||
r = r.rstrip(".,;:)»>")
|
||
r = re.sub(r"\.git$", "", r)
|
||
r = r.rstrip(".,;:")
|
||
return r or None
|
||
_KONZEPT_NAMES = ("KONZEPT.md", "Konzept.md", "konzept.md", "CONCEPT.md")
|
||
_konzept_cache: dict = {}
|
||
|
||
|
||
def _gitea_creds() -> tuple | None:
|
||
"""(user, token) fuer den oeffentlichen Gitea-Host aus ~/.git-credentials."""
|
||
try:
|
||
cred = Path("~/.git-credentials").expanduser()
|
||
for line in cred.read_text(encoding="utf-8").splitlines():
|
||
m = re.match(r"https://([^:]+):([^@]+)@" + re.escape(_GITEA_HOST), line.strip())
|
||
if m:
|
||
return m.group(1), m.group(2)
|
||
except Exception:
|
||
pass
|
||
return None
|
||
|
||
|
||
def _extract_repo(text: str) -> str | None:
|
||
"""Ziel-Repo (owner/name) aus dem Ergebnis-/Body-Text der Karte ziehen."""
|
||
if not text:
|
||
return None
|
||
m = _GITEA_REPO_RX.search(text)
|
||
if m:
|
||
return _clean_repo(m.group(1))
|
||
m = _REPO_FALLBACK_RX.search(text)
|
||
return _clean_repo(m.group(1)) if m else None
|
||
|
||
|
||
def konzept_of(task_id: str) -> dict:
|
||
"""Das KONZEPT.md des vorbereiteten IDE-Projekts holen (aus dem Gitea-Repo)."""
|
||
if not _available():
|
||
return {"available": False}
|
||
if not _TASK_ID_RX.match(task_id or ""):
|
||
return {"available": True, "error": "Ungültige Aufgaben-Nummer."}
|
||
if task_id in _konzept_cache:
|
||
return _konzept_cache[task_id]
|
||
|
||
# Repo aus Ergebnis + Body der Karte ziehen (show --json trägt beides).
|
||
try:
|
||
r = _hermes(["show", task_id, "--json"], timeout=30)
|
||
d = _parse_json(r.stdout)
|
||
except Exception:
|
||
return {"available": True, "error": "Karte nicht lesbar."}
|
||
# `kanban show --json` legt die Task-Felder verschachtelt unter "task" ab und
|
||
# liefert die jüngste Abschluss-Zusammenfassung als "latest_summary".
|
||
haystack, task_status = "", None
|
||
if isinstance(d, dict):
|
||
t = d.get("task") if isinstance(d.get("task"), dict) else d
|
||
task_status = t.get("status")
|
||
haystack = "\n".join(str(x or "") for x in (
|
||
t.get("result"), t.get("body"), d.get("latest_summary")))
|
||
for ev in d.get("events") or []:
|
||
if isinstance(ev, dict) and ev.get("kind") == "completed":
|
||
haystack += "\n" + str((ev.get("payload") or {}).get("summary") or "")
|
||
repo = _extract_repo(haystack)
|
||
if not repo:
|
||
return {"available": True, "repo": None,
|
||
"error": "Noch kein Repo hinterlegt — das Projekt wird gerade vorbereitet."}
|
||
|
||
creds = _gitea_creds()
|
||
auth = creds if creds else None
|
||
konzept, used_name = "", None
|
||
try:
|
||
import httpx
|
||
with httpx.Client(timeout=10.0, follow_redirects=True) as c:
|
||
for name in _KONZEPT_NAMES:
|
||
url = f"https://{_GITEA_HOST}/{repo}/raw/branch/main/{name}"
|
||
resp = c.get(url, auth=auth)
|
||
if resp.status_code == 200 and resp.text.strip():
|
||
konzept, used_name = resp.text, name
|
||
break
|
||
except Exception:
|
||
log.warning("ideen: KONZEPT-Abruf für %s (%s) fehlgeschlagen", task_id, repo, exc_info=True)
|
||
return {"available": True, "repo": repo, "error": "Konzept konnte nicht geladen werden."}
|
||
|
||
if not konzept:
|
||
return {"available": True, "repo": repo,
|
||
"error": "KONZEPT.md im Repo (noch) nicht gefunden."}
|
||
clone_url = f"https://{_GITEA_HOST}/{repo}.git"
|
||
data = {"available": True, "repo": repo, "clone_url": clone_url,
|
||
"datei": used_name, "konzept": konzept[:60000]}
|
||
# Fertige Konzepte ändern sich selten → cachen (aber nur, wenn die Karte done ist).
|
||
if task_status == "done":
|
||
if len(_konzept_cache) >= 20:
|
||
_konzept_cache.pop(next(iter(_konzept_cache)))
|
||
_konzept_cache[task_id] = data
|
||
return data
|
||
|
||
|
||
def log_of(task_id: str) -> dict:
|
||
"""Worker-Log einer Idee (hermes kanban log <id>), 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
|