feat: ide-lane hardening (ci ampel in auftragsbuch, hermes auto-bugfix, cleanup)
Ampel / ampel (push) Failing after 22s
Ampel / ampel (push) Failing after 22s
This commit is contained in:
@@ -26,6 +26,7 @@ import re
|
||||
import shutil
|
||||
import subprocess
|
||||
import time
|
||||
import urllib.request
|
||||
from pathlib import Path
|
||||
|
||||
from config import MODELS_DIR
|
||||
@@ -60,6 +61,46 @@ _KANDIDAT_RX = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._ -]{0,120}\.md$")
|
||||
|
||||
_fetch_cache: dict = {}
|
||||
_FETCH_EVERY = 30.0 # s — Gitea nicht bei jedem UI-Poll anfragen
|
||||
_CI_CACHE: dict = {} # { (repo, branch, head_sha): {"status": "success", "ts": time} }
|
||||
|
||||
def _gitea_creds() -> tuple | None:
|
||||
try:
|
||||
cred = Path("~/.git-credentials").expanduser()
|
||||
for line in cred.read_text(encoding="utf-8").splitlines():
|
||||
m = re.match(r"https://([^:]+):([^@]+)@git\.tobisniceshomelab", line.strip())
|
||||
if m: return m.group(1), m.group(2)
|
||||
except Exception:
|
||||
pass
|
||||
return None
|
||||
|
||||
def _fetch_ci_status(repo: str, branch: str, head_sha: str) -> str | None:
|
||||
# Nur auf der Box (wo creds liegen) sinnvoll
|
||||
creds = _gitea_creds()
|
||||
if not creds: return None
|
||||
user, token = creds
|
||||
cache_key = (repo, branch, head_sha)
|
||||
cached = _CI_CACHE.get(cache_key)
|
||||
if cached and (time.time() - cached["ts"] < 30 or cached["status"] in ("success", "failure", "skipped")):
|
||||
return cached["status"]
|
||||
|
||||
full_repo = "Hitonabi/mission-control-v2" if repo == "mc2" else "Hitonabi/lucy"
|
||||
url = f"http://192.168.178.153:3000/api/v1/repos/{full_repo}/actions/runs?branch={urllib.parse.quote(branch)}&limit=1"
|
||||
try:
|
||||
req = urllib.request.Request(url, headers={"Authorization": "token " + token})
|
||||
with urllib.request.urlopen(req, timeout=5) as r:
|
||||
data = json.load(r)
|
||||
runs = data if isinstance(data, list) else data.get("runs", data.get("workflow_runs", []))
|
||||
if runs:
|
||||
run = runs[0]
|
||||
if run.get("head_sha", "").startswith(head_sha[:7]):
|
||||
st = run.get("status")
|
||||
if st:
|
||||
_CI_CACHE[cache_key] = {"status": st, "ts": time.time()}
|
||||
return st
|
||||
except Exception as e:
|
||||
log.warning("auftragsbuch: CI-Status fetch fehler: %s", e)
|
||||
return None
|
||||
|
||||
|
||||
|
||||
def _available() -> bool:
|
||||
@@ -176,6 +217,10 @@ def _repo_items(repo: str, statuses: dict, gutachten: dict) -> list[dict]:
|
||||
stempel = gutachten.get(f"{repo}:{branch}")
|
||||
if not (isinstance(stempel, dict) and stempel.get("commit_ts") == ts_val):
|
||||
stempel = None
|
||||
|
||||
head_sha = (_git(repo, ["rev-parse", ref]).stdout or "").strip()
|
||||
ci_status = _fetch_ci_status(repo, branch, head_sha) if head_sha else None
|
||||
|
||||
stat = (_git(repo, ["diff", "--shortstat", f"origin/main...{ref}"]).stdout or "").strip()
|
||||
files_raw = (_git(repo, ["diff", "--name-status", f"origin/main...{ref}"]).stdout or "").splitlines()
|
||||
files = []
|
||||
@@ -209,6 +254,7 @@ def _repo_items(repo: str, statuses: dict, gutachten: dict) -> list[dict]:
|
||||
# es im Repo liegt. Lucy wird dagegen IMMER am PC gebaut (kein dist im Repo).
|
||||
"frontend_ohne_build": repo == "mc2" and frontend_src and not frontend_dist,
|
||||
"status": None if (status or {}).get("state") == "eingespielt" and ahead > 0 else status,
|
||||
"ci_status": ci_status,
|
||||
})
|
||||
except Exception:
|
||||
log.warning("auftragsbuch: Branch %s (%s) nicht lesbar", branch, repo, exc_info=True)
|
||||
|
||||
@@ -1,195 +0,0 @@
|
||||
"""
|
||||
Mini-Job-System: Hintergrund-Prozesse mit Live-Log + Download-Fortschritt.
|
||||
Portiert aus Mission Control v1 (jobengine.py). In-Memory, ein Daemon-Thread je Job.
|
||||
"""
|
||||
|
||||
import glob
|
||||
import os
|
||||
import shlex
|
||||
import subprocess
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
|
||||
JOBS: dict[str, dict] = {}
|
||||
_PROCS: dict[str, subprocess.Popen] = {}
|
||||
_LOG_CAP = 400
|
||||
|
||||
|
||||
def _append_log(job: dict, line: str) -> None:
|
||||
job["log"].append(line)
|
||||
if len(job["log"]) > _LOG_CAP:
|
||||
del job["log"][0]
|
||||
|
||||
|
||||
def _pump_output(job: dict, stream) -> None:
|
||||
"""Liest byteweise; `\\r` (tqdm/hf-Fortschritt) überschreibt die letzte Zeile."""
|
||||
buf = b""
|
||||
overwrite = False
|
||||
pending_cr = False
|
||||
|
||||
def commit():
|
||||
line = buf.decode("utf-8", "replace")
|
||||
if overwrite and job["log"]:
|
||||
job["log"][-1] = line
|
||||
else:
|
||||
_append_log(job, line)
|
||||
|
||||
while True:
|
||||
ch = stream.read(1)
|
||||
if not ch:
|
||||
break
|
||||
if pending_cr:
|
||||
pending_cr = False
|
||||
if ch == b"\n":
|
||||
commit(); overwrite = False; buf = b""
|
||||
continue
|
||||
commit(); overwrite = True; buf = b""
|
||||
if ch == b"\r":
|
||||
pending_cr = True
|
||||
elif ch == b"\n":
|
||||
commit(); overwrite = False; buf = b""
|
||||
else:
|
||||
buf += ch
|
||||
if pending_cr:
|
||||
commit(); overwrite = True; buf = b""
|
||||
if buf:
|
||||
commit()
|
||||
|
||||
|
||||
def _run_job(job_id: str, args: list[str], env: dict | None = None, sudo_password: str | None = None):
|
||||
job = JOBS[job_id]
|
||||
job["state"] = "running"
|
||||
try:
|
||||
actual_args = list(args)
|
||||
if sudo_password is not None:
|
||||
for i, arg in enumerate(actual_args):
|
||||
if isinstance(arg, str):
|
||||
actual_args[i] = arg.replace("sudo -n", "sudo -S").replace("sudo ", "sudo -S ")
|
||||
|
||||
proc = subprocess.Popen(
|
||||
actual_args, stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
|
||||
stdin=subprocess.PIPE if sudo_password is not None else None,
|
||||
bufsize=0,
|
||||
env={**os.environ, **(env or {})},
|
||||
)
|
||||
_PROCS[job_id] = proc
|
||||
|
||||
if sudo_password is not None and proc.stdin:
|
||||
proc.stdin.write((sudo_password + "\n").encode("utf-8"))
|
||||
proc.stdin.flush()
|
||||
proc.stdin.close()
|
||||
|
||||
_pump_output(job, proc.stdout)
|
||||
proc.wait()
|
||||
job["returncode"] = proc.returncode
|
||||
job["state"] = "canceled" if job.get("canceled") else ("done" if proc.returncode == 0 else "failed")
|
||||
|
||||
# Check if failed due to sudo authorization failure
|
||||
if proc.returncode != 0 and job["log"]:
|
||||
log_str = "\n".join(job["log"])
|
||||
if "a password is required" in log_str or "password" in log_str.lower() or "sudo:" in log_str:
|
||||
job["sudo_failed"] = True
|
||||
except Exception as exc:
|
||||
_append_log(job, f"[mc] Fehler: {exc}")
|
||||
job["state"] = "failed"
|
||||
job["returncode"] = -1
|
||||
finally:
|
||||
_PROCS.pop(job_id, None)
|
||||
job["finished_at"] = time.time()
|
||||
cb = job.pop("_on_done", None)
|
||||
if cb and job["state"] == "done":
|
||||
try:
|
||||
cb()
|
||||
except Exception as exc:
|
||||
_append_log(job, f"[mc] Nachbearbeitung-Fehler: {exc}")
|
||||
|
||||
|
||||
def attach_download_progress(job_id: str, local_dir: str, total_bytes: int) -> None:
|
||||
"""Fortschritt in % aus wachsenden *.incomplete-Dateien (hf schreibt sie)."""
|
||||
if not total_bytes or total_bytes <= 0:
|
||||
return
|
||||
job = JOBS.get(job_id)
|
||||
if job is not None:
|
||||
job["progress"] = 0
|
||||
job["total_bytes"] = total_bytes
|
||||
|
||||
def _watch():
|
||||
pat = os.path.join(local_dir, ".cache", "huggingface", "download", "**", "*.incomplete")
|
||||
prev_t = prev_b = None
|
||||
rate = 0.0
|
||||
while True:
|
||||
j = JOBS.get(job_id)
|
||||
if not j or j["state"] in ("done", "failed", "canceled"):
|
||||
break
|
||||
try:
|
||||
inc = glob.glob(pat, recursive=True)
|
||||
cur = sum(os.path.getsize(f) for f in inc) if inc else 0
|
||||
if cur:
|
||||
j["progress"] = min(99, int(cur * 100 / total_bytes))
|
||||
j["done_bytes"] = cur
|
||||
now = time.time()
|
||||
if prev_t is not None and now > prev_t and cur >= prev_b:
|
||||
inst = (cur - prev_b) / (now - prev_t)
|
||||
rate = inst if rate == 0 else 0.3 * inst + 0.7 * rate
|
||||
if rate > 0:
|
||||
j["rate_bps"] = rate
|
||||
j["eta_s"] = int((total_bytes - cur) / rate)
|
||||
prev_t, prev_b = now, cur
|
||||
except Exception:
|
||||
pass
|
||||
time.sleep(1.0)
|
||||
j = JOBS.get(job_id)
|
||||
if j and j["state"] == "done":
|
||||
j["progress"] = 100
|
||||
j.pop("eta_s", None)
|
||||
|
||||
threading.Thread(target=_watch, daemon=True).start()
|
||||
|
||||
|
||||
def start_job(args: list[str], label: str, env: dict | None = None, on_done=None,
|
||||
sudo_password: str | None = None, group: str | None = None) -> str:
|
||||
job_id = uuid.uuid4().hex[:12]
|
||||
# Mask password in log if present in args
|
||||
log_args = list(args)
|
||||
JOBS[job_id] = {
|
||||
"id": job_id, "label": label, "state": "queued", "group": group,
|
||||
"log": ["$ " + " ".join(shlex.quote(a) for a in log_args)],
|
||||
"returncode": None, "started_at": time.time(), "finished_at": None,
|
||||
}
|
||||
if on_done:
|
||||
JOBS[job_id]["_on_done"] = on_done
|
||||
threading.Thread(target=_run_job, args=(job_id, args, env, sudo_password), daemon=True).start()
|
||||
return job_id
|
||||
|
||||
|
||||
def active_in_group(group: str) -> dict | None:
|
||||
"""Erster laufender/wartender Job einer Gruppe (z.B. 'maintenance'), sonst None.
|
||||
Basis für den Wartungs-Riegel: nur EIN System-Update gleichzeitig."""
|
||||
for j in JOBS.values():
|
||||
if j.get("group") == group and j.get("state") in ("running", "queued"):
|
||||
return j
|
||||
return None
|
||||
|
||||
|
||||
def cancel_job(job_id: str) -> bool:
|
||||
job = JOBS.get(job_id)
|
||||
if not job or job["state"] in ("done", "failed", "canceled"):
|
||||
return False
|
||||
job["canceled"] = True
|
||||
_append_log(job, "[mc] Abbruch angefordert…")
|
||||
proc = _PROCS.get(job_id)
|
||||
if proc is not None:
|
||||
try:
|
||||
proc.terminate()
|
||||
except Exception:
|
||||
pass
|
||||
else:
|
||||
job["state"] = "canceled"
|
||||
job["finished_at"] = time.time()
|
||||
return True
|
||||
|
||||
|
||||
def public_jobs() -> list[dict]:
|
||||
"""Jobs ohne interne Felder (_on_done) für die API."""
|
||||
return [{k: v for k, v in j.items() if not k.startswith("_")} for j in JOBS.values()]
|
||||
@@ -1,190 +0,0 @@
|
||||
"""
|
||||
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 []
|
||||
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)
|
||||
@@ -1,178 +0,0 @@
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import psutil
|
||||
import time
|
||||
from datetime import datetime, timezone
|
||||
import httpx
|
||||
|
||||
from services import maintenance, announce, llamaswap, auftragsbuch, discover
|
||||
from config import V1_UPSTREAM
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# Wenn aktiv, wird der Report um die angegebene Stunde (0-23) lokaler Zeit verschickt
|
||||
SCHEDULE_HOUR = 8
|
||||
|
||||
async def fetch_tech_news() -> str:
|
||||
"""Holt die Top 3 Tech-News von HackerNews als Futter."""
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10.0) as client:
|
||||
top_ids_resp = await client.get("https://hacker-news.firebaseio.com/v0/topstories.json")
|
||||
top_ids = top_ids_resp.json()[:3]
|
||||
news_items = []
|
||||
for item_id in top_ids:
|
||||
item_resp = await client.get(f"https://hacker-news.firebaseio.com/v0/item/{item_id}.json")
|
||||
item = item_resp.json()
|
||||
title = item.get('title', 'Ohne Titel')
|
||||
url = item.get('url', f"https://news.ycombinator.com/item?id={item_id}")
|
||||
news_items.append(f"- {title} ({url})")
|
||||
return "\n".join(news_items) if news_items else "Keine News gefunden."
|
||||
except Exception as e:
|
||||
log.warning(f"Fehler beim Abrufen der News: {e}")
|
||||
return "News konnten nicht abgerufen werden."
|
||||
|
||||
async def generate_and_send_report() -> bool:
|
||||
"""Sammelt System-Status, fragt das lokale LLM nach einer Zusammenfassung als Lucy und sendet via Telegram."""
|
||||
log.info("Sysadmin-Report gestartet...")
|
||||
|
||||
# 1. Daten sammeln
|
||||
ram = psutil.virtual_memory()
|
||||
ram_gb_total = round(ram.total / (1024**3), 1)
|
||||
ram_gb_used = round(ram.used / (1024**3), 1)
|
||||
cpu_percent = psutil.cpu_percent(interval=1)
|
||||
disk = psutil.disk_usage('/')
|
||||
disk_gb_free = round(disk.free / (1024**3), 1)
|
||||
|
||||
try:
|
||||
upd = maintenance.updates()
|
||||
except Exception:
|
||||
upd = {"os": 0, "engine": False, "swap": False, "hermes": False}
|
||||
|
||||
running_models = llamaswap.get_running_models()
|
||||
|
||||
# News
|
||||
news_text = await fetch_tech_news()
|
||||
|
||||
# Discover: Bessere Modelle
|
||||
disc_data = discover.safe_discover(ram_gb_total)
|
||||
rec_models = []
|
||||
if disc_data and "categories" in disc_data:
|
||||
for c in disc_data["categories"]:
|
||||
if c.get("recommended"):
|
||||
rec_models.append(c["recommended"])
|
||||
rec_str = ", ".join(rec_models[:3]) if rec_models else "Keine neuen Empfehlungen"
|
||||
|
||||
# Auftragsbuch
|
||||
try:
|
||||
auftraege_data = auftragsbuch.list_proposals()
|
||||
open_count = auftraege_data.get("open_count", 0)
|
||||
auftraege_str = f"{open_count} offene Aufgaben (Karten), die auf dich warten." if open_count > 0 else "Das Auftragsbuch ist leer (Keine offenen Aufgaben)."
|
||||
except Exception as e:
|
||||
log.warning(f"Fehler im Sysadmin Report bei Auftragsbuch: {e}")
|
||||
auftraege_str = "Fehler beim Lesen des Auftragsbuchs."
|
||||
|
||||
events = announce.list_recent(limit=15)
|
||||
event_lines = []
|
||||
for e in events:
|
||||
ts_str = datetime.fromtimestamp(e["ts"]).strftime("%d.%m. %H:%M")
|
||||
text = e.get('text', '')
|
||||
if len(text) > 100:
|
||||
text = text[:100] + "..."
|
||||
event_lines.append(f"- [{ts_str}] {e.get('subject', '')}: {text}")
|
||||
|
||||
events_str = "\n".join(event_lines) if event_lines else "Keine besonderen Vorkommnisse."
|
||||
|
||||
evidence = f"""
|
||||
Hardware: RAM {ram_gb_used}/{ram_gb_total} GB ({ram.percent}%), CPU {cpu_percent}%, Speicher {disk_gb_free} GB frei.
|
||||
Laufende Modelle: {', '.join(running_models) if running_models else 'Keine'}
|
||||
Updates anstehend: OS: {upd.get('os', 0)}, Engine: {upd.get('engine', False)}, Swap: {upd.get('swap', False)}, Hermes: {upd.get('hermes', False)}
|
||||
Offene Freigaben für Commander (Auftragsbuch): {auftraege_str}
|
||||
Modell-Empfehlungen für diese Box: {rec_str}
|
||||
|
||||
Heutige Tech-News:
|
||||
{news_text}
|
||||
|
||||
Kürzliche Ereignisse (Chronik):
|
||||
{events_str}
|
||||
"""
|
||||
|
||||
prompt = f"""Du bist Lucy, die KI-Sysadmin der Mission-Control-Box (einem lokalen AI-Stack).
|
||||
Du KENNST die Box in- und auswendig. Dein Boss braucht einen gehaltvollen, klugen Morgen-Bericht via Telegram.
|
||||
|
||||
Regeln:
|
||||
1. Wie geht's der Box? (Kurzer Hardware Check, check ob Modelle laufen).
|
||||
2. News: Übersetze die Tech-News kurz auf Deutsch und übernimm IMMER die Quellen-URLs in deine Nachricht!
|
||||
3. Updates: Sind welche offen?
|
||||
4. Modelle & Aufgaben: Empfiehl nur Modelle, die nicht ohnehin schon laufen. Erinnere den Commander an offene Aufgaben im Auftragsbuch.
|
||||
5. Benutze Emojis, mach es lesbar (Bulletpoints).
|
||||
6. Chronik nur erwähnen, wenn es Probleme gab (sonst weglassen).
|
||||
7. Antworte DIREKT mit dem Text der Nachricht (keine Einleitung).
|
||||
8. SEI INFORMATIV, professionell-frech und zeige, dass du das System verstehst.
|
||||
|
||||
Rohe Daten:
|
||||
{evidence}
|
||||
"""
|
||||
|
||||
endpoint = "http://127.0.0.1:9001/v1/chat/completions" # Internes Gateway für model: auto
|
||||
|
||||
# Da das Gateway (Router) auf :9010 lauscht, können wir auch das nutzen,
|
||||
# aber Port 8080 (llama-swap engine) ist der direkteste Weg für den internen LLM-Call.
|
||||
model = "gpt-oss-120b" # Fallback, falls auto nicht geht, llama-swap routet das meist passend
|
||||
|
||||
req_body = {
|
||||
"model": "auto",
|
||||
"messages": [
|
||||
{"role": "system", "content": "Du bist Lucy, KI-Sysadmin der Box."},
|
||||
{"role": "user", "content": prompt}
|
||||
],
|
||||
"max_tokens": 1000,
|
||||
"temperature": 0.4
|
||||
}
|
||||
|
||||
digest = ""
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=180.0) as client:
|
||||
resp = await client.post(endpoint, json=req_body)
|
||||
resp.raise_for_status()
|
||||
data = resp.json()
|
||||
digest = data["choices"][0]["message"].get("content", "").strip()
|
||||
except Exception as e:
|
||||
log.error(f"Fehler beim LLM-Aufruf für Sysadmin-Report: {e}")
|
||||
# Fallback, falls LLM nicht erreichbar ist
|
||||
digest = f"🤖 [LLM offline] Hier sind die rohen Daten:\n{evidence}"
|
||||
|
||||
if digest:
|
||||
announce.notify_telegram("[🌅 Morgen-Digest]", digest)
|
||||
log.info("Sysadmin-Report via Telegram versendet.")
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
async def report_loop() -> None:
|
||||
"""Täglicher Background-Task für den Sysadmin-Report um SCHEDULE_HOUR Uhr."""
|
||||
log.info("Sysadmin-Report Background-Loop gestartet.")
|
||||
while True:
|
||||
try:
|
||||
now = datetime.now()
|
||||
# Finde die Zeit bis zum nächsten SCHEDULE_HOUR:00
|
||||
target_hour = SCHEDULE_HOUR
|
||||
if now.hour >= target_hour:
|
||||
# Nächster Tag
|
||||
seconds_until = ((24 - now.hour - 1) * 3600) + ((60 - now.minute - 1) * 60) + (60 - now.second) + (target_hour * 3600)
|
||||
else:
|
||||
# Selber Tag
|
||||
seconds_until = ((target_hour - now.hour - 1) * 3600) + ((60 - now.minute - 1) * 60) + (60 - now.second)
|
||||
|
||||
log.info(f"Sysadmin-Report: Warte {seconds_until} Sekunden bis zum nächsten Report.")
|
||||
await asyncio.sleep(seconds_until)
|
||||
|
||||
await generate_and_send_report()
|
||||
|
||||
# Warte kurz, um nicht sofort wieder zu triggern, falls es extrem schnell ging
|
||||
await asyncio.sleep(60)
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
log.error(f"Unerwarteter Fehler im Sysadmin-Report-Loop: {e}", exc_info=True)
|
||||
await asyncio.sleep(300)
|
||||
Reference in New Issue
Block a user