"""FastAPI-App: REST + WebSocket-Live-Board + statisches Frontend. Der State-Snapshot liefert IMMER alle Karten (kein 20er-Cap) — das Frontend scrollt je Spalte; Deltas kommen über /ws.""" import asyncio import logging import subprocess from pathlib import Path from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import FileResponse from fastapi.staticfiles import StaticFiles from pydantic import BaseModel import agents import db import guide import ledger import pipeline from config import FRONTEND_DIST, PROJECT_ROOT, topic_name_ok from ws import hub logging.basicConfig(level=logging.INFO, format="%(asctime)s %(name)s %(levelname)s %(message)s") log = logging.getLogger("creator2.main") app = FastAPI(title="creator2") app.add_middleware(CORSMiddleware, allow_origins=["http://localhost:5173"], allow_methods=["*"], allow_headers=["*"]) @app.on_event("startup") async def startup() -> None: db.connect() hub.bind_loop(asyncio.get_running_loop()) db.on_change = hub.push class TopicNeu(BaseModel): name: str titel: str = "" art: str = "thema" # thema | uni provider: str = "" class StartAuftrag(BaseModel): budget: int | None = None class UebenAntwort(BaseModel): artefakt_id: int richtig: bool class AutoFlag(BaseModel): ebene: str an: bool def _ebenen_info(t: dict, run: dict | None) -> list[dict]: """Je Ebene: fertig/laeuft/auto + Tokens (Läufe seit letztem Entfernen) + Dauer (letzter Lauf).""" flags = db.uj(t["auto"], {}) resets = db.uj(t["resets"], {}) out = [] for name, _, marke in pipeline.EBENEN: rows = db.query( "SELECT SUM(e.tok_in+e.tok_out) AS tok," " (julianday(MAX(e.ts))-julianday(MIN(e.ts)))*86400 AS dauer" " FROM events e JOIN runs r ON e.run_id=r.id" " WHERE r.topic=? AND e.ebene=? AND e.id>? GROUP BY e.run_id ORDER BY e.run_id", (t["name"], name, resets.get(name, 0))) out.append({ "name": name, "fertig": pipeline._fertig_ab(t["status"], marke), "laeuft": bool(run and run["status"] == "running" and run["ebene"] == name), "auto": bool(flags.get(name, True)), "tokens": sum(int(r["tok"] or 0) for r in rows), "dauer_s": int(rows[-1]["dauer"] or 0) if rows else 0, }) return out @app.get("/api/topics") def topics_liste(): out = [] for t in db.query("SELECT * FROM topics ORDER BY erstellt DESC"): run = db.one("SELECT * FROM runs WHERE topic=? ORDER BY id DESC", (t["name"],)) out.append({**t, "run": run}) return out @app.post("/api/topics") def topic_anlegen(auftrag: TopicNeu): if not auftrag.name.strip(): raise HTTPException(400, "Name fehlt") # streng validieren: der Name wird Pfadkomponente (rmtree, Snapshots) und # Shell-Argument (Transfer) — „..", „/" oder „;curl…" wären sonst gefährlich if not topic_name_ok(auftrag.name): raise HTTPException(400, "Name: nur Buchstaben/Ziffern/-/_ (Anfang alphanumerisch)") if db.one("SELECT name FROM topics WHERE name=?", (auftrag.name,)): raise HTTPException(409, "Topic existiert") if auftrag.art not in ("thema", "uni"): raise HTTPException(400, "art muss thema|uni sein") from config import DEFAULT_PROVIDER db.insert("topics", name=auftrag.name, titel=auftrag.titel or auftrag.name, art=auftrag.art, provider=auftrag.provider or DEFAULT_PROVIDER) return {"ok": True} @app.post("/api/topics/{topic}/start") async def lauf_start(topic: str, auftrag: StartAuftrag): # async: lauf_starten hängt den Lauf per ensure_future an DIESEN Event-Loop if not db.one("SELECT name FROM topics WHERE name=?", (topic,)): raise HTTPException(404, "unbekanntes Topic") try: run_id = pipeline.lauf_starten(topic, auftrag.budget) except RuntimeError as e: raise HTTPException(409, str(e)) return {"run_id": run_id} @app.post("/api/topics/{topic}/stop") async def lauf_stop(topic: str): pipeline.lauf_stoppen(topic) return {"ok": True} @app.get("/api/topics/{topic}/state") def state(topic: str): t = db.one("SELECT * FROM topics WHERE name=?", (topic,)) if not t: raise HTTPException(404, "unbekanntes Topic") run = db.one("SELECT * FROM runs WHERE topic=? ORDER BY id DESC", (topic,)) atome = db.query("SELECT * FROM atome WHERE topic=? ORDER BY id", (topic,)) fuer_atom = {} for a in db.query( "SELECT atom_id, SUM(status='verifiziert') AS ok, COUNT(*) AS n FROM artefakte" " WHERE atom_id IN (SELECT id FROM atome WHERE topic=?) GROUP BY atom_id", (topic,)): fuer_atom[a["atom_id"]] = {"verifiziert": a["ok"], "gesamt": a["n"]} return { "topic": t, "run": run, "ebenen": _ebenen_info(t, run), "verbraucht": ledger.verbraucht(run["id"]) if run else 0, "quellen": db.query("SELECT id, art, titel, url, runde, status, rolle FROM quellen" " WHERE topic=? ORDER BY id", (topic,)), # Kandidaten nur zeigen, solange der Konsens noch nichts bestätigt hat "soll": db.query("SELECT * FROM soll WHERE topic=? AND status='bestaetigt'" " ORDER BY id", (topic,)) or db.query("SELECT * FROM soll WHERE topic=? ORDER BY id", (topic,)), "atome": [{**a, "artefakte": fuer_atom.get(a["id"], {"verifiziert": 0, "gesamt": 0})} for a in atome], "lernziele": db.query("SELECT * FROM lernziele WHERE topic=? ORDER BY id", (topic,)), "bausteine": db.query( "SELECT b.*, s.stage FROM bausteine b LEFT JOIN sections s ON s.baustein_id=b.id" " WHERE b.topic=? ORDER BY b.ord", (topic,)), # Je Ebene nur der JÜNGSTE Messstand: offene Befunde älterer Läufe sind # veraltete Momentaufnahmen (aak: 1244 Anker-Leichen bei 0 echten) — der # letzte Lauf, der die Ebene gemessen hat, ist die Wahrheit. "befunde": db.query( "SELECT b.* FROM befunde b JOIN runs r ON r.id=b.run_id" " WHERE r.topic=? AND b.status='offen'" " AND b.run_id = (SELECT MAX(b2.run_id) FROM befunde b2" " JOIN runs r2 ON r2.id=b2.run_id" " WHERE r2.topic=? AND b2.ebene=b.ebene)" " ORDER BY b.id", (topic, topic)), "agenten": agents.aktive_agenten(), } @app.get("/api/health") def health(): return {"ok": True} # Themen-Transfer nur im lokalen Dev-Betrieb: der Server-Container (erkennbar an # /.dockerenv) hat weder ssh-Schlüssel noch das make-Gegenstück. _TRANSFER_LOKAL = not Path("/.dockerenv").exists() @app.get("/api/transfer") def transfer_info(): return {"verfuegbar": _TRANSFER_LOKAL} @app.post("/api/topics/{topic}/transfer/{richtung}") async def transfer_ausfuehren(topic: str, richtung: str): if not _TRANSFER_LOKAL: raise HTTPException(403, "Transfer nur lokal verfügbar") if richtung not in ("push", "pull"): raise HTTPException(404, "unbekannte Richtung") # Name gegen die Whitelist prüfen UND Existenz sichern — der Name fließt in # ein make-Shell-Rezept; ohne diese Schranke wäre er ein Injektionsvektor. if not topic_name_ok(topic) or not db.one("SELECT name FROM topics WHERE name=?", (topic,)): raise HTTPException(404, "unbekanntes Topic") res = await asyncio.to_thread( subprocess.run, ["make", f"server-{richtung}", f"TOPIC={topic}"], cwd=PROJECT_ROOT, capture_output=True, text=True, timeout=600) if res.returncode: raise HTTPException(500, (res.stderr or res.stdout)[-400:]) return {"ok": True} @app.get("/api/topics/{topic}/guide") def guide_holen(topic: str): # kein markdown-Feld: das Frontend rendert aus kapitel; Roh-Markdown # verdoppelte nur die Payload (aak: +0,43 MB ungenutzt) return {"kapitel": guide.kapitel_struktur(topic)} @app.post("/api/topics/{topic}/soll-reset") async def soll_reset(topic: str): if not db.one("SELECT name FROM topics WHERE name=?", (topic,)): raise HTTPException(404, "unbekanntes Topic") pipeline.soll_reset(topic) return {"ok": True} @app.post("/api/topics/{topic}/voll-reset") async def voll_reset(topic: str): if not db.one("SELECT name FROM topics WHERE name=?", (topic,)): raise HTTPException(404, "unbekanntes Topic") pipeline.voll_reset(topic) return {"ok": True} @app.patch("/api/topics/{topic}/auto") async def auto_setzen(topic: str, flag: AutoFlag): t = db.one("SELECT auto FROM topics WHERE name=?", (topic,)) if not t: raise HTTPException(404, "unbekanntes Topic") flags = db.uj(t["auto"], {}) flags[flag.ebene] = flag.an db.update("topics", "name", topic, auto=db.j(flags)) return {"ok": True} @app.post("/api/topics/{topic}/ebene/{ebene}/entfernen") async def ebene_entfernen(topic: str, ebene: str): if not db.one("SELECT name FROM topics WHERE name=?", (topic,)): raise HTTPException(404, "unbekanntes Topic") try: pipeline.ebenen_entfernen(topic, ebene) except ValueError as e: raise HTTPException(400, str(e)) return {"ok": True} @app.delete("/api/topics/{topic}") async def topic_loeschen(topic: str): if not db.one("SELECT name FROM topics WHERE name=?", (topic,)): raise HTTPException(404, "unbekanntes Topic") pipeline.topic_loeschen(topic) return {"ok": True} @app.get("/api/topics/{topic}/ueben") def ueben_faellig(topic: str): rows = db.query( "SELECT ar.id, ar.inhalt, a.titel, a.level, COALESCE(l.box,1) AS box," " COALESCE(l.faellig, date('now')) AS faellig" " FROM artefakte ar JOIN atome a ON ar.atom_id=a.id" " LEFT JOIN leitner l ON l.artefakt_id=ar.id" " WHERE a.topic=? AND ar.typ='flashcard' AND ar.status='verifiziert'" " AND COALESCE(l.faellig, date('now')) <= date('now') ORDER BY box, ar.id", (topic,)) return [{**r, "inhalt": db.uj(r["inhalt"], {})} for r in rows] _INTERVALLE = {1: 0, 2: 1, 3: 3, 4: 7, 5: 14} # Leitner: Tage bis zur Wiedervorlage @app.post("/api/ueben/antwort") def ueben_antwort(a: UebenAntwort): row = db.one("SELECT * FROM leitner WHERE artefakt_id=?", (a.artefakt_id,)) box = min(5, (row["box"] if row else 1) + 1) if a.richtig else 1 tage = _INTERVALLE[box] if row: db.execute("UPDATE leitner SET box=?, faellig=date('now', ?) WHERE artefakt_id=?", (box, f"+{tage} day", a.artefakt_id)) else: db.execute("INSERT INTO leitner(artefakt_id, box, faellig) VALUES(?,?,date('now', ?))", (a.artefakt_id, box, f"+{tage} day")) return {"box": box} @app.get("/api/topics/{topic}/baustein/{baustein_id}/karten") def baustein_karten(topic: str, baustein_id: int): # Flashcards genau dieses Bausteins — NICHT faellig-gefiltert (gerade gelesen → jetzt # prüfen). Beantwortet wird über /api/ueben/antwort, Leitner-Spacing bleibt intakt. rows = db.query( "SELECT ar.id, ar.inhalt, a.titel, a.level, COALESCE(l.box,1) AS box" " FROM artefakte ar JOIN atome a ON ar.atom_id=a.id" " LEFT JOIN leitner l ON l.artefakt_id=ar.id" " WHERE a.topic=? AND a.baustein_id=? AND ar.typ='flashcard'" " AND ar.status='verifiziert' ORDER BY a.ord, ar.id", (topic, baustein_id)) return [{**r, "inhalt": db.uj(r["inhalt"], {})} for r in rows] @app.get("/api/topics/{topic}/lernstand") def lernstand_holen(topic: str): return db.query("SELECT baustein_id, status, xp FROM lernstand WHERE topic=?", (topic,)) @app.post("/api/topics/{topic}/baustein/{baustein_id}/fertig") def baustein_fertig(topic: str, baustein_id: int): db.execute("INSERT INTO lernstand(topic, baustein_id, status) VALUES(?,?,'fertig')" " ON CONFLICT(topic, baustein_id) DO UPDATE SET status='fertig'", (topic, baustein_id)) return {"ok": True} @app.get("/api/runs/{run_id}/kennzahlen") def kennzahlen(run_id: int): return {"zeilen": ledger.kennzahlen(run_id), "verbraucht": ledger.verbraucht(run_id)} @app.websocket("/ws") async def websocket(ws: WebSocket): await hub.connect(ws) try: while True: await ws.receive_text() # Client sendet nichts Relevantes; hält die Verbindung except WebSocketDisconnect: hub.disconnect(ws) if FRONTEND_DIST.is_dir(): app.mount("/assets", StaticFiles(directory=FRONTEND_DIST / "assets"), name="assets") @app.get("/") def index(): return FileResponse(FRONTEND_DIST / "index.html")