Files
creator3/backend/main.py
2026-07-24 11:23:18 +02:00

238 lines
9.0 KiB
Python

"""FastAPI-App: Topics verwalten, Läufe starten/pausieren, Graph-Snapshot
(kompakt: Zähler statt Monolith), Task-Listen lazy, Guide, Kennzahlen."""
import asyncio
from contextlib import asynccontextmanager
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
from . import config, db, engine, graph, ledger, montage
from .ws import hub
@asynccontextmanager
async def lifespan(app: FastAPI):
db.connect()
engine.module_laden()
hub.bind_loop(asyncio.get_running_loop())
db.on_change = hub.push
yield
app = FastAPI(lifespan=lifespan)
app.add_middleware(CORSMiddleware, allow_origins=["http://localhost:5173"],
allow_methods=["*"], allow_headers=["*"])
class TopicNeu(BaseModel):
name: str
titel: str = ""
class StartAuftrag(BaseModel):
budget: int | None = None
class TopicEinstellung(BaseModel):
quellen_max: int | None = None # 0 = unbegrenzt, None = Default
@app.patch("/api/topics/{topic}")
async def topic_einstellen(topic: str, e: TopicEinstellung):
if not db.one("SELECT name FROM topics WHERE name=?", topic):
raise HTTPException(404)
if e.quellen_max is not None and (not 0 <= e.quellen_max <= 200
or e.quellen_max == 1):
# 1 Quelle kann nie bestätigen — Konsens braucht ≥2 unabhängige Belege
raise HTTPException(400, "quellen_max: 0 (unbegrenzt) oder 2-200")
db.update("topics", "name=?", (topic,), quellen_max=e.quellen_max)
return {"ok": True}
@app.get("/api/topics")
def topics_liste():
topics = db.query("SELECT * FROM topics ORDER BY erstellt DESC")
for t in topics:
t["run"] = db.one("SELECT id, status, grund, stufe FROM runs WHERE "
"topic=? ORDER BY id DESC LIMIT 1", t["name"])
return topics
@app.post("/api/topics")
def topic_anlegen(auftrag: TopicNeu):
if not config.TOPIC_NAME_RE.match(auftrag.name):
raise HTTPException(400, "ungültiger Name")
if db.one("SELECT name FROM topics WHERE name=?", auftrag.name):
raise HTTPException(409, "Topic existiert")
db.insert("topics", name=auftrag.name, titel=auftrag.titel)
return {"ok": True}
def _pipeline_daten_loeschen(topic: str) -> None:
"""Alle Pipeline-Daten eines Topics räumen (Snapshots inklusive)."""
for tabelle in db.PIPELINE_TABELLEN:
db.execute(f"DELETE FROM {tabelle} WHERE topic=?", topic)
db.execute("DELETE FROM befunde WHERE run_id IN "
"(SELECT id FROM runs WHERE topic=?)", topic)
db.execute("DELETE FROM sections WHERE baustein_id NOT IN "
"(SELECT id FROM bausteine)")
db.execute("DELETE FROM anker WHERE atom_id NOT IN (SELECT id FROM atome)")
ordner = (config.KORPUS_DIR / topic).resolve()
if ordner.is_dir() and ordner.parent == config.KORPUS_DIR.resolve():
import shutil
shutil.rmtree(ordner) # Guard: nur direkt unter KORPUS_DIR (Lektion!)
@app.delete("/api/topics/{topic}")
async def topic_loeschen(topic: str):
engine.lauf_stoppen(topic)
_pipeline_daten_loeschen(topic)
db.execute("DELETE FROM runs WHERE topic=?", topic)
db.execute("DELETE FROM topics WHERE name=?", topic)
return {"ok": True}
@app.post("/api/topics/{topic}/reset")
async def topic_reset(topic: str):
"""Pipeline auf null: Topic bleibt, Läufe/Ledger bleiben als Historie."""
if not db.one("SELECT name FROM topics WHERE name=?", topic):
raise HTTPException(404)
engine.lauf_stoppen(topic)
_pipeline_daten_loeschen(topic)
# Alte Läufe schließen — sonst resumed der nächste Start den alten Run
# samt altem Budget und Token-Zähler
db.execute("UPDATE runs SET status='stopped', beendet=? WHERE topic=? AND "
"status IN ('running','paused','budget')", db.now(), topic)
db.update("topics", "name=?", (topic,), status="neu")
return {"ok": True}
@app.post("/api/topics/{topic}/start")
async def start(topic: str, auftrag: StartAuftrag):
if not db.one("SELECT name FROM topics WHERE name=?", topic):
raise HTTPException(404)
try:
run_id = engine.lauf_starten(topic, auftrag.budget)
except RuntimeError as e:
raise HTTPException(409, str(e)) from e
return {"run_id": run_id}
@app.post("/api/topics/{topic}/pause")
async def pause(topic: str):
engine.lauf_pausieren(topic)
return {"ok": True}
@app.post("/api/topics/{topic}/stop")
async def stop(topic: str):
engine.lauf_stoppen(topic)
return {"ok": True}
@app.get("/api/topics/{topic}/graph")
def graph_snapshot(topic: str):
g = graph.get()
zaehler = {r["knoten"]: {} for r in []}
zaehler = {}
for r in db.query("SELECT knoten, status, COUNT(*) c FROM tasks WHERE "
"topic=? GROUP BY knoten, status", topic):
zaehler.setdefault(r["knoten"], {})[r["status"]] = r["c"]
gates = {r["knoten"]: {"status": r["status"], "runde": r["runde"]}
for r in db.query(
"SELECT knoten, status, MAX(runde) runde FROM gate_laeufe "
"WHERE topic=? GROUP BY knoten", topic)}
run = db.one("SELECT id, status, grund, stufe, budget_tokens FROM runs "
"WHERE topic=? ORDER BY id DESC LIMIT 1", topic)
if run:
run["verbraucht"] = ledger.verbraucht(run["id"])
return {"layout": g.layout(), "zaehler": zaehler, "gates": gates, "run": run}
@app.get("/api/topics/{topic}/knoten/{knoten}/tasks")
def knoten_tasks(topic: str, knoten: str):
tasks = db.query("SELECT id, item, status, art, runde, versuch, fehler, "
"gestartet, beendet FROM tasks WHERE topic=? AND knoten=? "
"ORDER BY id DESC LIMIT 200", topic, knoten)
befunde = db.query("SELECT art, item, detail, status FROM befunde WHERE "
"knoten=? AND run_id IN (SELECT id FROM runs WHERE topic=?) "
"ORDER BY id DESC LIMIT 100", knoten, topic)
return {"tasks": tasks, "befunde": befunde}
@app.get("/api/topics/{topic}/guide")
def guide_lesen(topic: str):
kapitel = []
for k in db.query("SELECT * FROM kapitel WHERE topic=? ORDER BY ord", topic):
sections = db.query(
"SELECT b.id, b.titel, s.text FROM bausteine b LEFT JOIN sections s "
"ON s.baustein_id=b.id WHERE b.kapitel_id=? ORDER BY b.ord", k["id"])
kapitel.append({"id": k["id"], "titel": k["titel"], "intro": k["intro"],
"sections": sections})
return {"kapitel": kapitel,
"fertig": (db.one("SELECT status FROM topics WHERE name=?", topic)
or {}).get("status") == "fertig"}
@app.get("/api/topics/{topic}/kennzahlen")
def kennzahlen(topic: str):
run = db.one("SELECT id, budget_tokens FROM runs WHERE topic=? "
"ORDER BY id DESC LIMIT 1", topic)
if not run:
return {"zeilen": [], "verbraucht": 0, "kosten_usd": 0, "coverage": None}
kz = ledger.kennzahlen(run["id"])
gesamt = db.one("SELECT COUNT(*) c FROM soll WHERE topic=? AND "
"status='bestaetigt'", topic)["c"]
gedeckt = db.one(
"SELECT COUNT(*) c FROM soll s WHERE s.topic=? AND s.status='bestaetigt' "
"AND (s.freispruch!='' OR EXISTS (SELECT 1 FROM atome a WHERE "
"a.soll_id=s.id AND a.status='aktiv' AND a.baustein_id IS NOT NULL))",
topic)["c"]
redundanz_offen = db.one(
"SELECT COUNT(*) c FROM befunde WHERE art='redundanz' AND status='offen' "
"AND run_id IN (SELECT id FROM runs WHERE topic=?)", topic)["c"]
eskaliert = db.one(
"SELECT COUNT(*) c FROM befunde WHERE status='eskaliert' AND run_id IN "
"(SELECT id FROM runs WHERE topic=?)", topic)["c"]
atome_aktiv = db.one("SELECT COUNT(*) c FROM atome WHERE topic=? AND "
"status='aktiv'", topic)["c"]
verdichtet = db.one("SELECT COUNT(*) c FROM soll WHERE topic=? AND "
"status='bestaetigt' AND verdichtet=1", topic)["c"]
from . import belege as belege_mod
kz["beleg_abdeckung"] = belege_mod.abdeckung_prozent(topic)
kz.update({"coverage": round(gedeckt / gesamt * 100, 1) if gesamt else None,
"soll_gesamt": gesamt, "soll_gedeckt": gedeckt,
"redundanz_offen": redundanz_offen, "eskaliert": eskaliert,
"atome_aktiv": atome_aktiv, "verdichtet_punkte": verdichtet,
"budget": run["budget_tokens"]})
return kz
@app.get("/api/health")
def health():
return {"ok": True}
@app.websocket("/ws")
async def ws_endpoint(ws: WebSocket):
await ws.accept()
hub.connect(ws)
try:
while True:
await ws.receive_text()
except WebSocketDisconnect:
hub.disconnect(ws)
_DIST = config.ROOT / "frontend" / "dist"
if _DIST.exists():
app.mount("/assets", StaticFiles(directory=_DIST / "assets"), name="assets")
@app.get("/")
def index():
return FileResponse(_DIST / "index.html")