Files
creator2/backend/main.py
2026-07-10 17:14:17 +02:00

291 lines
10 KiB
Python

"""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
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")
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,)),
"befunde": db.query(
"SELECT * FROM befunde WHERE run_id=? AND status='offen' ORDER BY id",
(run["id"],)) if run else [],
"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")
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/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")