diff --git a/deploy/compose.yml b/deploy/compose.yml new file mode 100644 index 0000000..afb7f4d --- /dev/null +++ b/deploy/compose.yml @@ -0,0 +1,67 @@ +# km -- der Wissensdienst als compose-Dienst (spine/core#8, #13). +# +# Zwei Container auf einem Volume: +# km API + Takt: zieht die Ablage, liest Gitea ein, indexiert +# (SQLite FTS5 + sqlite-vec), liefert alles Geaenderte ein +# silverbullet Oberflaeche auf demselben Verzeichnis -- lesen und aendern; +# km committet und pusht die Aenderungen im naechsten Takt +# +# Die Wahrheit ist das Git-Remote (KM_GIT_REMOTE); /srv/km/data/ablage ist +# ein Checkout, /srv/km/data/km-index.sqlite ein Cache. Beides ist ohne die +# Container lesbar bzw. neu baubar (Nachfolgefall). +# +# Secrets stehen in `.env` daneben (root-only) und NICHT hier: +# KM_GIT_USER / KM_GIT_PASSWORD Push an das Ablage-Repo +# KM_GITEA_USER / KM_GITEA_PASSWORD Lesen der Forge +# SB_USER user:passwort fuer die Oberflaeche +# +# Betrieb ueber `oci` (42i-oci): /srv/km/compose.yml, `oci enable km`. + +services: + km: + image: git.42i.org/spine/km:latest + build: + context: ../image + container_name: km + restart: unless-stopped + network_mode: host + environment: + KM_NAMESPACE: ${KM_NAMESPACE:-familie} + KM_CLASSIFICATION: ${KM_CLASSIFICATION:-home} + KM_GIT_REMOTE: ${KM_GIT_REMOTE} + KM_GIT_BRANCH: ${KM_GIT_BRANCH:-main} + KM_GIT_USER: ${KM_GIT_USER} + KM_GIT_PASSWORD: ${KM_GIT_PASSWORD} + KM_GIT_NAME: ${KM_GIT_NAME:-km} + KM_GIT_EMAIL: ${KM_GIT_EMAIL:-km@localhost} + KM_GITEA_URL: ${KM_GITEA_URL} + KM_GITEA_USER: ${KM_GITEA_USER} + KM_GITEA_PASSWORD: ${KM_GITEA_PASSWORD} + KM_GITEA_ORGS: ${KM_GITEA_ORGS:-} + KM_GITEA_REPOS: ${KM_GITEA_REPOS:-} + KM_INTERVAL: ${KM_INTERVAL:-300} + # Vektoren erst, wenn es einen Abnehmer gibt (Andreas 04.09.); die + # Ollama-Instanz auf s18 ist vom LXC aus nur mit VLAN5-Route erreichbar. + KM_EMBED: ${KM_EMBED:-0} + KM_OLLAMA_URL: ${KM_OLLAMA_URL:-http://10.18.5.30:11434} + KM_PORT: "8390" + volumes: + - /srv/km/data:/data + + silverbullet: + image: ghcr.io/silverbulletmd/silverbullet:v2 + container_name: km-ui + restart: unless-stopped + network_mode: host + environment: + SB_FOLDER: /space + SB_PORT: "3000" + # hive.home teilt den Pfadraum an der Wurzel (hive/core#61); die + # Oberflaeche liegt deshalb unter /km, die API unter /km/api. + SB_URL_PREFIX: /km + SB_USER: ${SB_USER} + # die Dateien gehoeren dem km-Container (root); ohne PUID raet + # SilverBullet die UID aus dem Ordner -- das ist dann dieselbe + PUID: "0" + volumes: + - /srv/km/data/ablage:/space diff --git a/deploy/health b/deploy/health new file mode 100755 index 0000000..1d5e47e --- /dev/null +++ b/deploy/health @@ -0,0 +1,7 @@ +#!/bin/sh +# 42i-hc (30-compose) ruft /srv//health: 0 = gesund. Prueft die API, +# nicht nur den Prozess -- ein Dienst, der antwortet, aber dessen letzter +# Lauf rot ist, meldet das ueber "ok": false. +out=$(curl -fsS --max-time 5 http://127.0.0.1:8390/health) || { echo "km: API antwortet nicht"; exit 1; } +printf '%s' "$out" | grep -q '"ok": *true' || { echo "km: letzter Lauf mit Fehler: $out"; exit 1; } +echo "km: ok" diff --git a/image/Dockerfile b/image/Dockerfile new file mode 100644 index 0000000..123b723 --- /dev/null +++ b/image/Dockerfile @@ -0,0 +1,31 @@ +# km -- der Wissensdienst (spine/core#8). Benannt nach der Faehigkeit. +# +# Ein Prozess: HTTP-API + Takt (ziehen, einlesen, indexieren, einliefern). +# Die Ablage ist ein Git-Checkout auf dem Volume /data/ablage, der Index eine +# SQLite-Datei daneben. SilverBullet als Oberflaeche ist ein eigener Container +# auf demselben Volume (deploy/compose.yml), nicht Teil dieses Images. +FROM docker.io/library/python:3.13-slim + +RUN apt-get update \ + && apt-get install -y --no-install-recommends git ca-certificates curl \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /app +COPY requirements.txt /app/ +RUN pip install --no-cache-dir -r requirements.txt + +COPY km/ /app/km/ +COPY entrypoint.sh /app/entrypoint.sh +RUN chmod +x /app/entrypoint.sh /app/km/*.py + +ENV KM_ABLAGE=/data/ablage \ + KM_PORT=8390 \ + PYTHONUNBUFFERED=1 + +VOLUME ["/data"] +EXPOSE 8390 + +HEALTHCHECK --interval=60s --timeout=5s --start-period=30s \ + CMD curl -fsS http://127.0.0.1:8390/health >/dev/null || exit 1 + +ENTRYPOINT ["/app/entrypoint.sh"] diff --git a/image/entrypoint.sh b/image/entrypoint.sh new file mode 100644 index 0000000..443d67e --- /dev/null +++ b/image/entrypoint.sh @@ -0,0 +1,6 @@ +#!/bin/sh +# Start, der laut scheitert: fehlt Notwendiges, haelt main.py mit einer Liste an. +set -eu +git config --global --add safe.directory '*' +git config --global pull.rebase true +exec python3 /app/km/main.py diff --git a/image/km/config.py b/image/km/config.py new file mode 100644 index 0000000..7cbd41c --- /dev/null +++ b/image/km/config.py @@ -0,0 +1,87 @@ +"""Konfiguration des km-Dienstes aus der Umgebung. + +Alles, was je Instanz verschieden ist, kommt von aussen (compose/Chart); +Secrets nur ueber Umgebungsvariablen, nie in Dateien im Image. Fehlt etwas +Notwendiges, haelt der Start mit einer Liste an -- ein Dienst, der startet +und nichts tut, ist teurer als einer, der gar nicht startet (spine/images). +""" + +from __future__ import annotations + +import os +from dataclasses import dataclass, field +from pathlib import Path + + +def _list(v: str) -> list[str]: + return [x.strip() for x in v.replace(";", ",").split(",") if x.strip()] + + +@dataclass +class Config: + ablage: Path + namespace: str + classification: str + git_remote: str | None + git_branch: str + git_user: str | None + git_password: str | None + git_name: str + git_email: str + gitea_url: str | None + gitea_user: str | None + gitea_password: str | None + gitea_orgs: list[str] = field(default_factory=list) + gitea_repos: list[str] = field(default_factory=list) + interval: int = 300 + embed: bool = False + ollama_url: str = "http://10.18.5.30:11434" + embed_model: str = "bge-m3" + listen: str = "0.0.0.0" + port: int = 8390 + + @property + def db(self) -> Path: + return self.ablage.parent / "km-index.sqlite" + + @property + def state(self) -> Path: + return self.ablage.parent / "km-ingest-state.json" + + @classmethod + def from_env(cls) -> "Config": + e = os.environ.get + cfg = cls( + ablage=Path(e("KM_ABLAGE", "/data/ablage")), + namespace=e("KM_NAMESPACE", ""), + classification=e("KM_CLASSIFICATION", "lan"), + git_remote=e("KM_GIT_REMOTE") or None, + git_branch=e("KM_GIT_BRANCH", "main"), + git_user=e("KM_GIT_USER") or None, + git_password=e("KM_GIT_PASSWORD") or None, + git_name=e("KM_GIT_NAME", "km"), + git_email=e("KM_GIT_EMAIL", "km@localhost"), + gitea_url=(e("KM_GITEA_URL") or "").rstrip("/") or None, + gitea_user=e("KM_GITEA_USER") or None, + gitea_password=e("KM_GITEA_PASSWORD") or None, + gitea_orgs=_list(e("KM_GITEA_ORGS", "")), + gitea_repos=_list(e("KM_GITEA_REPOS", "")), + interval=int(e("KM_INTERVAL", "300")), + embed=e("KM_EMBED", "0").lower() in ("1", "true", "yes", "ja"), + ollama_url=e("KM_OLLAMA_URL", "http://10.18.5.30:11434"), + embed_model=e("KM_EMBED_MODEL", "bge-m3"), + listen=e("KM_LISTEN", "0.0.0.0"), + port=int(e("KM_PORT", "8390")), + ) + missing = [] + if not cfg.namespace: + missing.append("KM_NAMESPACE (z. B. familie oder 42i)") + if cfg.git_remote and cfg.git_remote.startswith("http") and not (cfg.git_user and cfg.git_password): + missing.append("KM_GIT_USER und KM_GIT_PASSWORD (Push an KM_GIT_REMOTE)") + if cfg.gitea_url and not (cfg.gitea_user and cfg.gitea_password): + missing.append("KM_GITEA_USER und KM_GITEA_PASSWORD (Ingest von KM_GITEA_URL)") + if cfg.gitea_url and not (cfg.gitea_orgs or cfg.gitea_repos): + missing.append("KM_GITEA_ORGS oder KM_GITEA_REPOS") + if missing: + raise SystemExit("km: Konfiguration unvollstaendig:\n - " + "\n - ".join(missing)) + return cfg diff --git a/image/km/gitrepo.py b/image/km/gitrepo.py new file mode 100644 index 0000000..bf7e80f --- /dev/null +++ b/image/km/gitrepo.py @@ -0,0 +1,116 @@ +"""Die Ablage als Git-Checkout: klonen, ziehen, alles Geaenderte einliefern. + +Weg 1 aus spine/core#8: das Volume traegt ein Checkout, die Wahrheit liegt im +Remote, der Dienst committet und pusht selbst. Muster wie +info/ai/wissens_mcp/repo.py -- bei abgelehntem Push einmal rebasen und +wiederholen. Ein gescheiterter Push wird nicht verschluckt, sondern als +Zustand gemeldet (`status()`), damit niemand glaubt, die Wahrheit laege im +Remote, wenn sie nur lokal liegt. + +Zugangsdaten ueber eine ~/.netrc (0600), die der Start aus der Umgebung +schreibt -- nie in der URL, nie auf der Kommandozeile. +""" + +from __future__ import annotations + +import re +import subprocess +import threading +import time +from pathlib import Path +from urllib.parse import urlparse + +from config import Config + +_REJECTED = re.compile(r"\b(rejected|non-fast-forward|fetch first|behind)\b", re.IGNORECASE) + + +class Ablage: + def __init__(self, cfg: Config) -> None: + self.cfg = cfg + self.lock = threading.Lock() + self.last_push_error: str | None = None + self.last_push: float | None = None + + # ------------------------------------------------------------ Hilfen + def _git(self, *args: str, check: bool = True) -> subprocess.CompletedProcess: + return subprocess.run(["git", *args], cwd=self.cfg.ablage, check=check, + capture_output=True, text=True) + + def write_netrc(self) -> None: + if not (self.cfg.git_remote and self.cfg.git_user and self.cfg.git_password): + return + host = urlparse(self.cfg.git_remote).hostname + if not host: + return + netrc = Path.home() / ".netrc" + line = f"machine {host} login {self.cfg.git_user} password {self.cfg.git_password}\n" + existing = netrc.read_text() if netrc.exists() else "" + if f"machine {host} " not in existing: + netrc.write_text(existing + line) + netrc.chmod(0o600) + + # ------------------------------------------------------------ Lebenszyklus + def ensure(self) -> None: + """Checkout vorhanden machen: klonen, sonst init -- nie stillschweigend leer.""" + a = self.cfg.ablage + if (a / ".git").exists(): + return + a.parent.mkdir(parents=True, exist_ok=True) + if self.cfg.git_remote: + r = subprocess.run(["git", "clone", "--branch", self.cfg.git_branch, self.cfg.git_remote, str(a)], + capture_output=True, text=True) + if r.returncode != 0: + # leeres Remote: Branch gibt es noch nicht -> init und ersten Push + if "Remote branch" in r.stderr or "not found" in r.stderr or "warning: You appear" in r.stderr: + a.mkdir(parents=True, exist_ok=True) + self._git("init", "-b", self.cfg.git_branch) + self._git("remote", "add", "origin", self.cfg.git_remote) + else: + raise SystemExit(f"km: clone fehlgeschlagen: {r.stderr.strip()}") + else: + a.mkdir(parents=True, exist_ok=True) + self._git("init", "-b", self.cfg.git_branch) + self._git("config", "user.name", self.cfg.git_name) + self._git("config", "user.email", self.cfg.git_email) + + def pull(self) -> None: + if not self.cfg.git_remote: + return + r = self._git("pull", "--rebase", "--autostash", "origin", self.cfg.git_branch, check=False) + if r.returncode != 0 and "couldn't find remote ref" not in r.stderr: + print(f"km: pull: {r.stderr.strip()[:300]}", flush=True) + + def commit_all(self, message: str) -> bool: + """Alles Geaenderte einliefern (auch Aenderungen aus SilverBullet). True, wenn ein Commit entstand.""" + self._git("add", "-A") + if self._git("diff", "--cached", "--quiet", check=False).returncode == 0: + return False + self._git("-c", f"user.name={self.cfg.git_name}", "-c", f"user.email={self.cfg.git_email}", + "commit", "-q", "-m", message) + return True + + def push(self) -> None: + if not self.cfg.git_remote: + return + for _ in range(3): + r = self._git("push", "-u", "origin", self.cfg.git_branch, check=False) + if r.returncode == 0: + self.last_push_error = None + self.last_push = time.time() + return + if not _REJECTED.search(r.stderr): + self.last_push_error = r.stderr.strip()[:300] + print(f"km: push fehlgeschlagen: {self.last_push_error}", flush=True) + return + self._git("pull", "--rebase", "origin", self.cfg.git_branch, check=False) + self.last_push_error = "push nach 3 Versuchen abgelehnt" + print(f"km: {self.last_push_error}", flush=True) + + def head(self) -> str | None: + r = self._git("rev-parse", "--short", "HEAD", check=False) + return r.stdout.strip() or None + + def status(self) -> dict: + return {"head": self.head(), "remote": bool(self.cfg.git_remote), + "letzter_push": self.last_push, "push_fehler": self.last_push_error} diff --git a/image/km/ingest_gitea.py b/image/km/ingest_gitea.py index f9a6e28..6aebae6 100755 --- a/image/km/ingest_gitea.py +++ b/image/km/ingest_gitea.py @@ -62,6 +62,11 @@ class Forge: def host(self) -> str: return re.sub(r"^https?://", "", self.url).split("/")[0] + def login_with(self, user: str, password: str) -> None: + """Zugangsdaten direkt (im Dienst: aus der Umgebung, nie geloggt).""" + self.session.auth = (user, password) + self.session.headers["Accept"] = "application/json" + def login(self, secret_bin: str) -> None: def get(fld: str) -> str: r = subprocess.run( diff --git a/image/km/main.py b/image/km/main.py new file mode 100644 index 0000000..dc86f10 --- /dev/null +++ b/image/km/main.py @@ -0,0 +1,29 @@ +#!/usr/bin/env python3 +"""Einstieg des km-Dienstes: Konfiguration pruefen, Ablage sichern, Takt und API starten.""" + +from __future__ import annotations + +import sys +import threading + +from config import Config +from gitrepo import Ablage +from runner import Runner +from server import Api, serve + + +def main() -> None: + cfg = Config.from_env() + ablage = Ablage(cfg) + ablage.write_netrc() + ablage.ensure() + runner = Runner(cfg, ablage) + print(f"km: namespace={cfg.namespace} ablage={cfg.ablage} remote={'ja' if cfg.git_remote else 'nein'} " + f"gitea={cfg.gitea_url or '-'} takt={cfg.interval}s embed={'an' if cfg.embed else 'aus'}", + file=sys.stderr, flush=True) + threading.Thread(target=runner.loop, name="takt", daemon=True).start() + serve(cfg, Api(cfg, ablage, runner)) + + +if __name__ == "__main__": + main() diff --git a/image/km/runner.py b/image/km/runner.py new file mode 100644 index 0000000..8247b6a --- /dev/null +++ b/image/km/runner.py @@ -0,0 +1,100 @@ +"""Der Takt des Dienstes: ziehen, einlesen, indexieren, einliefern. + +Ein Durchlauf (`run_once`) ist idempotent und darf jederzeit wiederholt +werden. Er wird vom Timer (KM_INTERVAL) und von `POST /api/v1/refresh` +ausgeloest; spaeter vom Weckruf am Bus (spine/core#12). Zwei Durchlaeufe +laufen nie gleichzeitig -- der zweite wartet. +""" + +from __future__ import annotations + +import sys +import threading +import time +import traceback + +import index as km_index +from config import Config +from gitrepo import Ablage +from ingest_gitea import Forge, ingest_repo, load_state, save_state + + +class Runner: + def __init__(self, cfg: Config, ablage: Ablage) -> None: + self.cfg = cfg + self.ablage = ablage + self.lock = threading.Lock() + self.last: dict = {"zeit": None, "dauer": None, "fehler": None, "ergebnis": None} + self._wake = threading.Event() + + def wake(self) -> None: + self._wake.set() + + def run_once(self, full: bool = False) -> dict: + with self.lock: + t0 = time.time() + result: dict = {} + try: + self.ablage.pull() + result["ingest"] = self._ingest(full) + result["index"] = self._index() + if self.ablage.commit_all(self._message(result)): + self.ablage.push() + result["commit"] = self.ablage.head() + self.last = {"zeit": t0, "dauer": round(time.time() - t0, 1), "fehler": None, "ergebnis": result} + except Exception as e: # noqa: BLE001 -- der Takt darf nicht sterben + traceback.print_exc() + self.last = {"zeit": t0, "dauer": round(time.time() - t0, 1), "fehler": str(e)[:300], "ergebnis": result} + print(f"km: lauf {self.last['dauer']}s {self.last['ergebnis']} fehler={self.last['fehler']}", + file=sys.stderr, flush=True) + return self.last + + def _ingest(self, full: bool) -> dict | None: + c = self.cfg + if not c.gitea_url: + return None + forge = Forge(url=c.gitea_url, namespace=c.namespace, classification=c.classification) + forge.login_with(c.gitea_user or "", c.gitea_password or "") + repos = list(c.gitea_repos) + for org in c.gitea_orgs: + repos += forge.org_repos(org) + state = load_state(c.state) + seen = written = 0 + for repo in dict.fromkeys(repos): + s, w = ingest_repo(forge, repo, c.ablage, state, full) + seen += s + written += w + save_state(c.state, state) + return {"repos": len(repos), "gesehen": seen, "geschrieben": written} + + def _index(self) -> dict: + c = self.cfg + con = km_index.connect(c.db) + try: + files = sorted(p for p in c.ablage.rglob("*.md") if not any(s.startswith(".") for s in p.parts)) + changed = sum(km_index.index_file(con, c.ablage, f, c.namespace, 600, 60) for f in files) + gone = km_index.remove_missing(con, c.ablage) + out = {"dateien": len(files), "geaendert": changed, "entfernt": gone} + if c.embed and changed: + out["eingebettet"] = km_index.embed_pending( + con, km_index.Embedder(c.ollama_url, c.embed_model), None) + return out + finally: + con.close() + + @staticmethod + def _message(result: dict) -> str: + ing = result.get("ingest") or {} + idx = result.get("index") or {} + parts = [] + if ing.get("geschrieben"): + parts.append(f"{ing['geschrieben']} Vorgänge aus Gitea") + if idx.get("geaendert"): + parts.append(f"{idx['geaendert']} Dateien geändert") + return "km: " + (", ".join(parts) if parts else "Änderungen aus der Oberfläche") + + def loop(self) -> None: + while True: + self.run_once() + self._wake.wait(self.cfg.interval) + self._wake.clear() diff --git a/image/km/server.py b/image/km/server.py new file mode 100644 index 0000000..12c360d --- /dev/null +++ b/image/km/server.py @@ -0,0 +1,214 @@ +"""HTTP-API von km: read, search, summary, health, refresh. + +Bewusst Standardbibliothek (ThreadingHTTPServer): keine Abhaengigkeit, die +beim Bau des Images oder in einer Sitzung wegfliegen kann. JSON rein, JSON +raus; Fehler als {"fehler": ...} mit passendem Status. + + GET /health + GET /api/v1/search?q=…&mode=hybrid|fts|vec&limit=8&ns=&projekt=&typ=&zustand=&classification=a,b + GET /api/v1/read?path= + GET /api/v1/summary?ns=&projekt=&tage=7 + POST /api/v1/refresh[?full=1] -- Lauf anstossen (ziehen, einlesen, indexieren, einliefern) +""" + +from __future__ import annotations + +import json +import re +import sqlite3 +import sys +import threading +import time +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from urllib.parse import parse_qs, urlparse + +import index as km_index +from config import Config +from gitrepo import Ablage +from runner import Runner +from search import fts_query, rrf + + +class _Args: + """Filterobjekt in der Form, die search.py erwartet.""" + + def __init__(self, q: dict[str, list[str]]) -> None: + g = lambda k: (q.get(k) or [None])[0] # noqa: E731 + self.query = g("q") or "" + self.ns = g("ns") + self.projekt = g("projekt") + self.typ = g("typ") + self.zustand = g("zustand") + cl = g("classification") + self.classification = [c for c in cl.split(",") if c] if cl else None + self.mode = g("mode") or "hybrid" + self.limit = max(1, min(int(g("limit") or 8), 50)) + + +class Api: + def __init__(self, cfg: Config, ablage: Ablage, runner: Runner) -> None: + self.cfg = cfg + self.ablage = ablage + self.runner = runner + self.embedder = km_index.Embedder(cfg.ollama_url, cfg.embed_model) + + def _con(self) -> sqlite3.Connection: + return km_index.connect(self.cfg.db) + + # ------------------------------------------------------------ Endpunkte + def health(self) -> dict: + con = self._con() + try: + docs = con.execute("SELECT count(*) FROM docs").fetchone()[0] + embedded = con.execute("SELECT count(*) FROM docs WHERE eingebettet=1").fetchone()[0] + finally: + con.close() + last = self.runner.last + return {"ok": last.get("fehler") is None, "namespace": self.cfg.namespace, + "dokumente": docs, "eingebettet": embedded, "git": self.ablage.status(), + "letzter_lauf": last} + + def search(self, a: _Args) -> list[dict]: + from search import search_fts, search_vec + if not a.query.strip(): + raise ValueError("q fehlt") + con = self._con() + try: + k = a.limit * 3 + fts = search_fts(con, a, k) if a.mode in ("hybrid", "fts") else [] + vec = [] + if a.mode in ("hybrid", "vec"): + try: + vec = search_vec(con, a, k, self.embedder) + except Exception as e: # noqa: BLE001 -- ohne Embedding bleibt der Volltext + print(f"km: vektorsuche nicht verfuegbar: {e}", file=sys.stderr, flush=True) + if a.mode == "vec": + raise + if a.mode == "fts": + ranked = [(cid, -s) for cid, s in fts] + elif a.mode == "vec": + ranked = [(cid, -d) for cid, d in vec] + else: + ranked = sorted(rrf(fts, vec).items(), key=lambda x: -x[1]) + out, seen = [], set() + for cid, score in ranked: + row = con.execute( + "SELECT c.path, c.idx, c.text, d.titel, d.typ, d.zustand, d.aktualisiert, d.classification " + "FROM chunks c JOIN docs d ON d.path=c.path WHERE c.id=?", (cid,)).fetchone() + if not row or row[0] in seen: + continue + seen.add(row[0]) + out.append({"path": row[0], "titel": row[3], "typ": row[4], "zustand": row[5], + "aktualisiert": row[6], "classification": row[7], "score": round(score, 4), + "anriss": re.sub(r"\s+", " ", row[2])[:240]}) + if len(out) >= a.limit: + break + return out + finally: + con.close() + + def read(self, rel: str) -> dict: + if not rel or rel.startswith("/") or ".." in rel.split("/"): + raise ValueError("path ungueltig") + f = (self.cfg.ablage / rel).resolve() + if not str(f).startswith(str(self.cfg.ablage.resolve())) or not f.is_file(): + raise FileNotFoundError(rel) + text = f.read_text(encoding="utf-8") + meta, body = km_index.parse_frontmatter(text) + con = self._con() + try: + hin = [r[0] for r in con.execute("SELECT nach FROM verweise WHERE von=?", (rel,))] + her = [r[0] for r in con.execute("SELECT von FROM verweise WHERE nach=?", (self._as_ref(meta),))] if meta else [] + finally: + con.close() + return {"path": rel, "meta": meta, "text": text, "verweise": hin, "verwiesen_von": her} + + @staticmethod + def _as_ref(meta: dict) -> str: + return f"{meta.get('repo')}#{meta.get('nummer')}" if meta.get("repo") else "" + + def summary(self, q: dict[str, list[str]]) -> dict: + """Deterministischer Ueberblick; die verdichtende Fassung per Modell kommt spaeter (#8).""" + ns = (q.get("ns") or [self.cfg.namespace])[0] + projekt = (q.get("projekt") or [None])[0] + tage = int((q.get("tage") or ["7"])[0]) + seit = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(time.time() - tage * 86400)) + con = self._con() + try: + w, p = "d.namespace=?", [ns] + if projekt: + w += " AND d.projekt=?" + p.append(projekt) + counts = {r[0] or "-": r[1] for r in con.execute( + f"SELECT d.typ || '/' || coalesce(d.zustand,'-'), count(*) FROM docs d WHERE {w} GROUP BY 1", p)} + recent = [{"path": r[0], "titel": r[1], "typ": r[2], "zustand": r[3], "aktualisiert": r[4]} + for r in con.execute( + f"SELECT path, titel, typ, zustand, aktualisiert FROM docs d WHERE {w} AND aktualisiert >= ? " + "ORDER BY aktualisiert DESC LIMIT 50", [*p, seit])] + offen = [{"path": r[0], "titel": r[1], "aktualisiert": r[2]} for r in con.execute( + f"SELECT path, titel, aktualisiert FROM docs d WHERE {w} AND typ='vorgang' AND zustand='open' " + "ORDER BY aktualisiert DESC LIMIT 100", p)] + finally: + con.close() + return {"namespace": ns, "projekt": projekt, "seit": seit, "bestand": counts, + "geaendert": recent, "offene_vorgaenge": len(offen), "offen": offen} + + +def make_handler(api: Api): + class Handler(BaseHTTPRequestHandler): + server_version = "km/0.1" + + def _send(self, code: int, obj) -> None: + data = json.dumps(obj, ensure_ascii=False).encode("utf-8") + self.send_response(code) + self.send_header("Content-Type", "application/json; charset=utf-8") + self.send_header("Content-Length", str(len(data))) + self.end_headers() + self.wfile.write(data) + + def log_message(self, fmt, *args): # ruhig, ausser Fehler + if args and str(args[1]).startswith(("4", "5")): + sys.stderr.write("km: %s %s\n" % (self.address_string(), fmt % args)) + + def do_GET(self) -> None: + u = urlparse(self.path) + q = parse_qs(u.query) + try: + if u.path == "/health": + self._send(200, api.health()) + elif u.path == "/api/v1/search": + self._send(200, {"treffer": api.search(_Args(q))}) + elif u.path == "/api/v1/read": + self._send(200, api.read((q.get("path") or [""])[0])) + elif u.path == "/api/v1/summary": + self._send(200, api.summary(q)) + else: + self._send(404, {"fehler": "unbekannt", "pfade": ["/health", "/api/v1/search", "/api/v1/read", "/api/v1/summary", "POST /api/v1/refresh"]}) + except FileNotFoundError as e: + self._send(404, {"fehler": f"nicht gefunden: {e}"}) + except ValueError as e: + self._send(400, {"fehler": str(e)}) + except Exception as e: # noqa: BLE001 + self._send(500, {"fehler": str(e)[:300]}) + + def do_POST(self) -> None: + u = urlparse(self.path) + q = parse_qs(u.query) + if u.path == "/api/v1/refresh": + full = (q.get("full") or ["0"])[0] in ("1", "true") + if api.runner.lock.locked(): + self._send(202, {"status": "laeuft bereits"}) + return + threading.Thread(target=api.runner.run_once, args=(full,), daemon=True).start() + self._send(202, {"status": "angestossen", "full": full}) + else: + self._send(404, {"fehler": "unbekannt"}) + + return Handler + + +def serve(cfg: Config, api: Api) -> None: + httpd = ThreadingHTTPServer((cfg.listen, cfg.port), make_handler(api)) + print(f"km: API auf http://{cfg.listen}:{cfg.port}", file=sys.stderr, flush=True) + httpd.serve_forever()