Mache km zum Dienst: HTTP-API, Takt, Git-Ablage, Image und compose
Ein Prozess (image/km/main.py) mit API (read, search, summary, health, refresh) auf Standardbibliothek und einem Takt, der die Ablage zieht, Gitea einliest, indexiert und alles Geänderte committet und pusht — auch Änderungen aus SilverBullet. Die Ablage ist ein Git-Checkout auf dem Volume, der Index eine SQLite-Datei daneben; ein Start ohne nötige Umgebung hält mit einer Liste an. deploy/compose.yml: km plus SilverBullet auf demselben Volume, für den oci-Betrieb in hive.home unter /km (SB_URL_PREFIX). deploy/health für 42i-hc. Lokal gegen ein Bare-Remote und spine/core durchgespielt: Ingest, Index, Commit, Push, Suche, Read mit Kanten in beide Richtungen, und eine von Hand angelegte Seite landet im nächsten Takt im Remote. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
@@ -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
|
||||
Executable
+7
@@ -0,0 +1,7 @@
|
||||
#!/bin/sh
|
||||
# 42i-hc (30-compose) ruft /srv/<svc>/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"
|
||||
@@ -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"]
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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}
|
||||
@@ -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(
|
||||
|
||||
@@ -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()
|
||||
@@ -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()
|
||||
@@ -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=<relativer Pfad in der Ablage>
|
||||
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()
|
||||
Reference in New Issue
Block a user