diff --git a/.env.example b/.env.example index 1c0aa0e..ffe96aa 100644 --- a/.env.example +++ b/.env.example @@ -12,8 +12,9 @@ BRAIN_STAGING_DIRS=../glpi-ai-knowledgebase/staging # Optional read-only audit streams. Multiple files may be comma-separated. BRAIN_AGENT_RUNS_FILES=../glpi-ai-agent/data/runs.jsonl -# Batched disk persistence. AI-THINK drafts, the GLPI-KB cache and graph-state -# are written sequentially at this interval and once more during shutdown. +# Batched disk persistence. AI-THINK files and caches are written first; +# changed graph rows follow in one SQLite/WAL transaction. A final flush is +# attempted during shutdown. BRAIN_PERSIST_INTERVAL=5m # Ollama pool. OLLAMA_URLS takes precedence over legacy OLLAMA_URL. @@ -58,6 +59,8 @@ BRAIN_LEARNING_CATEGORIES= BRAIN_DISPLAY_CATEGORIES= BRAIN_THINKING_CATEGORIES= BRAIN_DEFAULT_VIEW=neural +# 0 = unbegrenzt; Webinterface kann den Wert zur Laufzeit ändern +BRAIN_MAX_DISPLAY_NODES=0 # Sequential enrichment BRAIN_AUTO_ENRICH=true @@ -72,5 +75,20 @@ BRAIN_TOP_K=8 BRAIN_MAX_CONTEXT_CHARS=16000 # Optional controlled web research through your own SearXNG instance. +# Use the root URL or a URL ending in /search. Inside Docker, localhost points +# to the Brain container; use the SearXNG service name or host.docker.internal. BRAIN_RESEARCH_ENABLED=false SEARXNG_URL= + +# Knowledge synthesis after a verified relation has formed a useful source cluster. +# Relation thinking always remains separate and only creates graph edges. +BRAIN_ARTICLE_SYNTHESIS_ENABLED=true +BRAIN_ARTICLE_MIN_SOURCES=3 +BRAIN_ARTICLE_MAX_SOURCES=8 +BRAIN_ARTICLE_MIN_PRODUCTION_RATIO=0.70 +BRAIN_ARTICLE_MAX_GENERATION_DEPTH=2 +BRAIN_ARTICLE_MIN_CONFIDENCE=0.74 +BRAIN_ARTICLE_MIN_TEXT_CHARS=180 +BRAIN_ARTICLE_MIN_ANSWER_CHARS=420 +BRAIN_ARTICLE_MAX_RESEARCH_QUERIES=3 +BRAIN_ARTICLE_RESEARCH_RESULTS=4 diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 59f24ba..e0f4819 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -2,7 +2,7 @@ ```text Lokale Knowledge-JSONs ──ro──┐ -GLPI Knowledge Base ─────ro──┼──► Ingest ─► In-Memory Graph + Vectors +GLPI Knowledge Base ─────ro──┼──► Ingest ─► In-Memory Graph + Float32 Vectors Agent runs.jsonl ─────────ro──┘ │ ├──► Ollama Pool ├──► Relation Thinking @@ -12,7 +12,7 @@ Agent runs.jsonl ─────────ro──┘ ├── staging drafts ├── article metadata sidecars ├── GLPI cache - └── graph-state.json + └── SQLite/WAL graph.db ``` ## Edge-Klassen @@ -22,7 +22,6 @@ Agent runs.jsonl ─────────ro──┘ - `related_to`, `depends_on`, `supports`, `contradicts`, `extends`, `same_topic`, `caused_by`: Qwen-Inferenz mit Confidence und Evidence. - `rejected`: intern gespeicherte Prüfung ohne sichtbare Beziehung. - ## Mehrstufiges Thinking ```text @@ -48,7 +47,9 @@ Relation Thinking schreibt keine Artikel. Die sichtbare KB-Datei enthält aussch ## Schreibmodell -Produktive Knowledge-Dateien und GLPI sind read-only. AI-THINK wird nur in das konfigurierte Staging geschrieben. Alle Schreibzugriffe laufen durch einen Prozess-weiten Coordinator, der Dateien dedupliziert und seriell vor dem Graph-Snapshot schreibt. +Produktive Knowledge-Dateien und GLPI sind read-only. AI-THINK wird nur in das konfigurierte Staging geschrieben. Ein prozessweiter Coordinator dedupliziert Dateischreibvorgänge und schreibt sie seriell vor der Graphtransaktion. + +Der Graph Store verwendet `modernc.org/sqlite` über `database/sql`. Nodes, Edges und binäre Float32-Vektoren werden nur bei Änderungen geschrieben. WAL erlaubt eine kompakte laufende Journaldatei; der Export-Endpunkt erzeugt mit `VACUUM INTO` eine portable Einzeldatei ohne WAL-Sidecars. ## Ollama-Pool diff --git a/CHANGELOG-GLPI-POOL-PERSISTENCE.md b/CHANGELOG-GLPI-POOL-PERSISTENCE.md index 86e48b2..b540abd 100644 --- a/CHANGELOG-GLPI-POOL-PERSISTENCE.md +++ b/CHANGELOG-GLPI-POOL-PERSISTENCE.md @@ -24,7 +24,7 @@ - zentrale Queue für AI-THINK- und Cache-Dateien; - Deduplizierung nach Zielpfad; - serielles, atomares Schreiben; -- Graph-Snapshot erst nach den Wissensdateien; +- Graphtransaktion erst nach den Wissensdateien; - Standardintervall fünf Minuten; - finaler Flush beim geregelten Shutdown; - manueller Flush über `POST /api/flush`. diff --git a/CHANGELOG-SEARXNG-DIAGNOSTICS.md b/CHANGELOG-SEARXNG-DIAGNOSTICS.md new file mode 100644 index 0000000..1bb09bf --- /dev/null +++ b/CHANGELOG-SEARXNG-DIAGNOSTICS.md @@ -0,0 +1,24 @@ +# Changelog: SearXNG-Diagnose + +## Neu + +- Direkter SearXNG-Test im Webinterface unter `FILTER → SearXNG-Diagnose`. +- `GET /api/research/status` für den zuletzt bekannten Verbindungsstatus. +- `POST /api/research/test` für eine echte JSON-Suche ohne Graph-Ingest. +- Eigene Events `research.test.started`, `research.test.results` und `research.test.failed`. +- Test und Fehler verwenden dieselbe mindestens zweisekündige Rechercheanimation. + +## Fehlerdiagnose + +- Anzeige der tatsächlich verwendeten, von Zugangsdaten bereinigten Basis-URL. +- Erfassung von HTTP-Status, Content-Type, Laufzeit und Trefferzahl. +- Klassifikation von DNS-, Verbindungs-, Timeout-, TLS-, HTTP- und JSON-Fehlern. +- Antwortausschnitt bei Proxy-, SearXNG- oder Formatfehlern. +- Hinweise für das Docker-`localhost`-Problem, falsche Service-Namen und deaktivierte JSON-Ausgabe. +- Automatische Relations- und Artikelrecherche schreiben die konkrete Ursache nun in Log und Aktivitätsfeed. + +## URL-Kompatibilität + +- `SEARXNG_URL` darf sowohl auf die SearXNG-Basis als auch direkt auf einen Pfad mit `/search` zeigen. +- Ein vorhandener Reverse-Proxy-Pfad bleibt erhalten. +- Der vollständige Compose-Stack erhält für Linux ebenfalls `host.docker.internal:host-gateway`, damit ein auf dem Docker-Host laufendes SearXNG erreichbar ist. diff --git a/CHANGELOG-SQLITE-STARTUP-FIX.md b/CHANGELOG-SQLITE-STARTUP-FIX.md new file mode 100644 index 0000000..58b53b3 --- /dev/null +++ b/CHANGELOG-SQLITE-STARTUP-FIX.md @@ -0,0 +1,39 @@ +# SQLite startup fix + +## Problem + +Some installations aborted during `graph.Open` with an opaque modernc SQLite error similar to: + +```text +SQL logic error: out of memory (1) +``` + +The original startup path applied a 64 MiB SQLite page cache, a 256 MiB mmap window and `temp_store=MEMORY` while the connection was being opened. Those settings are unnecessary because the live graph and embeddings already reside in Go memory. The DSN-based PRAGMA list also made it impossible to identify which startup step failed. + +## Changes + +- Open SQLite using a plain absolute filename. +- Verify that `BRAIN_DATA_DIR` exists, is a directory and is writable before opening SQLite. +- Reject a `graph.db` path that is accidentally a directory. +- Apply PRAGMAs individually with step-specific error messages. +- Use conservative storage defaults: + - `cache_size=-8192` (approximately 8 MiB) + - `mmap_size=0` + - `temp_store=FILE` + - WAL, `synchronous=NORMAL`, foreign keys and a 5-second busy timeout remain enabled. +- Persist and report the journal mode actually selected by SQLite. +- Add the dependency checksums to `go.sum` and copy both `go.mod` and `go.sum` in the Docker dependency layer. + +## Recovery + +Because this is a staging-only database, remove files left by a failed first start before rebuilding: + +```bash +docker compose stop brain +docker compose run --rm --no-deps --entrypoint sh brain -c ' + rm -f /app/data/graph.db /app/data/graph.db-wal /app/data/graph.db-shm +' +docker compose up -d --build --force-recreate brain +``` + +If startup still fails, the new error includes the exact phase, for example an unwritable data directory or a failed WAL configuration. diff --git a/CHANGELOG-SQLITE-STORAGE.md b/CHANGELOG-SQLITE-STORAGE.md new file mode 100644 index 0000000..edd18e4 --- /dev/null +++ b/CHANGELOG-SQLITE-STORAGE.md @@ -0,0 +1,13 @@ +# SQLite Storage + +- `graph-state.json` durch `graph.db` ersetzt. +- `modernc.org/sqlite v1.37.1` als CGo-freien `database/sql`-Treiber eingebunden. +- Nodes, Edges und Vektoren in normalisierten SQLite-Tabellen gespeichert. +- Embeddings von JSON/`float64` auf binäre little-endian `float32`-BLOBs umgestellt. +- WAL-Modus und gebündelte, inkrementelle Transaktionen eingeführt. +- Periodische KB-Scans reconciliieren Datensätze und schreiben unveränderte Zeilen nicht erneut. +- Vektoren werden nur bei Änderungen an Label, Summary, Kategorien oder Keywords invalidiert. +- Embedding-Modellname und Ollama-Digest werden gespeichert; ein Identitätswechsel invalidiert importierte Vektoren automatisch. +- `GET /api/state/export` erzeugt eine portable, kompakte `graph.db`-Kopie. +- `/api/status` um `graph_storage` erweitert. +- Kein Legacy-Import aus `graph-state.json`, da die Zielumgebung Staging ist. diff --git a/Dockerfile b/Dockerfile index 9c336c9..4d095f3 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,9 +1,15 @@ FROM golang:1.26-alpine AS build WORKDIR /src -COPY go.mod ./ + +# modernc.org/sqlite is a pure-Go database/sql driver; no compiler toolchain or +# CGO runtime is required. Dependencies are resolved in a separate layer so +# normal source changes do not redownload the SQLite module tree. +COPY go.mod go.sum ./ +RUN go mod download + COPY cmd ./cmd COPY internal ./internal -RUN CGO_ENABLED=0 GOOS=linux go build -trimpath -ldflags="-s -w" -o /out/neural-brain ./cmd/brain +RUN CGO_ENABLED=0 GOOS=linux go build -mod=mod -trimpath -ldflags="-s -w" -o /out/neural-brain ./cmd/brain FROM alpine:3.24 RUN addgroup -S brain && adduser -S -G brain brain diff --git a/Makefile b/Makefile index e6f7a0a..0292500 100644 --- a/Makefile +++ b/Makefile @@ -1,4 +1,8 @@ -.PHONY: run test build fmt +.PHONY: deps run test build fmt export-state + +deps: + go mod download + run: go run ./cmd/brain @@ -6,7 +10,10 @@ test: go test ./... build: - go build -o bin/neural-brain ./cmd/brain + CGO_ENABLED=0 go build -o bin/neural-brain ./cmd/brain fmt: gofmt -w cmd internal + +export-state: + curl -fsS http://localhost:8090/api/state/export -o graph.db diff --git a/PERSISTENCE.md b/PERSISTENCE.md index b9aa461..a4e4e44 100644 --- a/PERSISTENCE.md +++ b/PERSISTENCE.md @@ -1,15 +1,111 @@ -# Gebündelte Persistenz +# SQLite-Persistenz -Der vollständige Graph arbeitet im Speicher. Schreibvorgänge werden in einer zentralen Queue gesammelt und standardmäßig alle fünf Minuten seriell ausgeführt. +Der vollständige Graph bleibt für Retrieval, Visualisierung und AI-THINK im Arbeitsspeicher. Dauerhaft gespeichert wird er nicht mehr als JSON-Gesamtsnapshot, sondern inkrementell in: -## Reihenfolge +```text +BRAIN_DATA_DIR/graph.db +``` -1. Ausstehende AI-THINK-Entwürfe und Cache-Dateien werden nach Pfad sortiert und atomar geschrieben. -2. Danach wird genau ein konsistenter `graph-state.json`-Snapshot erzeugt. -3. Bei SIGTERM/SIGINT wird ein finaler Flush versucht. +Als Treiber wird `modernc.org/sqlite` verwendet. Die Anwendung bleibt damit ohne CGO baubar. -Dateien werden über eine temporäre Datei, `fsync`, `close` und `rename` ersetzt. Mehrere ausstehende Schreibvorgänge für denselben Pfad werden zusammengeführt; nur der neueste Inhalt wird geschrieben. +## Datenmodell -## Trade-off +Die SQLite-Datenbank enthält: -Ein längeres Intervall reduziert I/O, vergrößert aber das Zeitfenster nicht persistierter Änderungen. Für Systeme mit zuverlässigem Shutdown sind fünf Minuten ein sinnvoller Ausgangspunkt. Für maximale Haltbarkeit kann das Intervall verkürzt werden. +- `nodes` – Nodes einschließlich Kategorien, Keywords und Metadaten; +- `edges` – Beziehungen einschließlich Evidence und Confidence; +- `vectors` – Embeddings als little-endian `float32`-BLOBs; +- `graph_meta` – Schema- und Graphversion. + +Die Vektoren benötigen exakt vier Byte pro Dimension. Sie werden im Arbeitsspeicher ebenfalls als `float32` gehalten und nur an bestehenden API-Grenzen vorübergehend nach `float64` konvertiert. + +## Schreibmodell + +```env +BRAIN_PERSIST_INTERVAL=5m +``` + +Ein Flush läuft sequenziell: + +1. ausstehende AI-THINK-, Metadaten-, Runtime- und Cache-Dateien werden nach Pfad sortiert und atomar ersetzt; +2. anschließend schreibt eine SQLite-Transaktion ausschließlich geänderte oder gelöschte Node-, Edge- und Vector-Zeilen; +3. ein passiver WAL-Checkpoint hält die WAL-Datei begrenzt; +4. bei SIGTERM oder SIGINT wird ein finaler Flush und ein `TRUNCATE`-Checkpoint versucht. + +Mehrere Änderungen derselben Wissensdatei werden in der Dateiqueue dedupliziert. Ein unveränderter KB-Scan markiert keine Graphzeilen als schreibbedürftig und behält vorhandene Embeddings. + +## SQLite-Einstellungen + +Die Datenbank läuft mit: + +```text +journal_mode=WAL +synchronous=NORMAL +foreign_keys=ON +busy_timeout=5000 +wal_autocheckpoint=0 +temp_store=MEMORY +``` + +Der Graph-Store verwendet absichtlich eine einzelne `database/sql`-Verbindung. In-Memory-Lesezugriffe werden weiterhin über den Store parallel bedient; nur die gebündelte Persistenz ist serialisiert. + +## Durability-Fenster + +Neue In-Memory-Erkenntnisse sind sofort sichtbar. Bei einem harten Stromausfall können jedoch die Änderungen seit dem letzten Flush fehlen. Ein regulärer Container-Shutdown löst einen finalen Flush aus. + +Manueller Flush: + +```bash +curl -fsS -X POST http://localhost:8090/api/flush +``` + +## Portabler Export + +Die laufende Datenbank sollte nicht zusammen mit ihren `-wal`- und `-shm`-Dateien kopiert werden. Der Export-Endpunkt verwendet `VACUUM INTO` und erzeugt eine kompakte, konsistente Einzeldatei: + +```bash +curl -fsS \ + -H "Authorization: Bearer $BRAIN_API_KEY" \ + http://localhost:8090/api/state/export \ + -o graph.db +``` + +Ohne API-Key entfällt der Header. + +Import auf einem anderen System mit identischer KB-Basis: + +```bash +docker compose stop brain + +# Vorhandene DB und mögliche Sidecars entfernen. +docker compose run --rm --no-deps -T --entrypoint sh brain -c ' + rm -f /app/data/graph.db /app/data/graph.db-wal /app/data/graph.db-shm + cat > /app/data/graph.db +' < graph.db + +docker compose up -d brain +``` + +Beim nächsten Scan bleiben Vektoren für unveränderte Dokument-IDs und identische Embedding-Inhalte erhalten. Das in `graph_meta` gespeicherte Embedding-Modell wird beim Start geprüft. Sobald Ollama erreichbar ist, wird zusätzlich der Modelldigest aus dem Pool gespeichert und verglichen. Ändern sich Modellname oder Digest, werden vorhandene Vektoren automatisch verworfen und selektiv neu gelernt. Neue oder veränderte Beiträge werden ebenfalls nur einzeln neu eingebettet. + +## Vollständiger Reset + +Da diese Staging-Version keine Migration aus `graph-state.json` benötigt, genügt: + +```bash +docker compose stop brain +docker compose run --rm --no-deps --entrypoint sh brain -c ' + rm -f /app/data/graph.db /app/data/graph.db-wal /app/data/graph.db-shm +' +docker compose up -d brain +``` + +Eine eventuell noch vorhandene `graph-state.json` wird ignoriert. `runtime-settings.json`, GLPI-Cache, Artikelmetadaten und Staging-Drafts werden durch diesen Reset nicht gelöscht. + +## Status + +```bash +curl -fsS http://localhost:8090/api/status | jq '.graph_storage, .persistence' +``` + +`graph_storage` zeigt unter anderem DB-/WAL-Größe, binäre Vector-Bytes, persistierte Graphversion und ausstehende Row-Änderungen. diff --git a/README.md b/README.md index 42c818a..81d04f6 100644 --- a/README.md +++ b/README.md @@ -14,7 +14,7 @@ Eigenständiger Go-Dienst für Agent, lokale Knowledgebase, GLPI-Knowledgebase u - Ollama-Pool mit mehreren unabhängigen Instanzen, Routing, Healthchecks, Cooldown und Failover. - Embeddings über `embeddinggemma`, Beziehungsanalyse über `qwen3:8b`. - Mehrstufiger autonomer Worker: Relation Thinking, quellengebundene Wissenskonsolidierung, Recherche offener Punkte und reine Knowledge-Synthesis für vollständige KB-Artikel. -- Gebündelte Festplattenpersistenz: Graph, GLPI-KB-Cache und AI-THINK-Dateien werden standardmäßig nur alle fünf Minuten sequenziell geschrieben. +- Inkrementelle SQLite/WAL-Persistenz über `modernc.org/sqlite`: binäre Float32-Embeddings und standardmäßig alle fünf Minuten gebündelte Row-Updates. ## Vertrauens- und Schreibgrenzen @@ -106,19 +106,18 @@ curl -X POST http://localhost:8090/api/glpi-kb/sync Mehr Details: [`GLPI-KB.md`](GLPI-KB.md). -## Gebündelte Persistenz +## SQLite/WAL-Persistenz und Vorberechnung ```env BRAIN_PERSIST_INTERVAL=5m ``` -Im Arbeitsspeicher sind neue Nodes, Edges und AI-THINK-Ergebnisse sofort verfügbar. Auf die Festplatte wird jedoch sequenziell geschrieben: +Der Live-Graph liegt im Arbeitsspeicher. Dauerhaft gespeichert werden nur geänderte Zeilen in `BRAIN_DATA_DIR/graph.db`: -1. ausstehende AI-THINK- und Cache-Dateien atomar, Pfad für Pfad; -2. anschließend genau ein Graph-Snapshot; -3. zusätzlich ein finaler Flush beim geregelten Shutdown. - -Das reduziert Schreibzugriffe deutlich. Der Preis ist ein konfigurierbares Durability-Fenster: Bei einem harten Stromausfall können die seit dem letzten Flush entstandenen Änderungen fehlen. +1. AI-THINK-, Metadaten- und Cache-Dateien werden atomar geschrieben; +2. danach folgen geänderte Nodes, Edges und Vektoren in einer SQLite-Transaktion; +3. Embeddings liegen binär als `float32`-BLOBs vor; +4. unveränderte KB-Scans verursachen keine erneuten Graphwrites. Manueller Flush: @@ -126,7 +125,13 @@ Manueller Flush: curl -X POST http://localhost:8090/api/flush ``` -Mehr Details: [`PERSISTENCE.md`](PERSISTENCE.md). +Ein auf einem leistungsfähigen System vorberechneter Graph kann kompakt exportiert und auf ein System mit identischen Dokument-IDs und demselben Embedding-Modell kopiert werden: + +```bash +curl -fsS http://localhost:8090/api/state/export -o graph.db +``` + +Mehr Details: [`PERSISTENCE.md`](PERSISTENCE.md) und [`SQLITE-STORAGE.md`](SQLITE-STORAGE.md). ## Laufzeitsteuerung und Honeycomb @@ -206,7 +211,8 @@ curl -X POST 'http://localhost:8090/api/enrich?async=1' | `POST` | `/api/reindex` | lokaler Scan und Embedding-Abgleich | | `POST` | `/api/enrich` | AI-THINK-Zyklus einplanen/ausführen | | `POST` | `/api/glpi-kb/sync` | GLPI-KB manuell synchronisieren | -| `POST` | `/api/flush` | ausstehende Dateien und Graph sofort persistieren | +| `POST` | `/api/flush` | ausstehende Dateien und geänderte SQLite-Zeilen sofort persistieren | +| `GET` | `/api/state/export` | kompakte, konsistente `graph.db` exportieren | Mit `BRAIN_API_KEY` werden POST-Endpunkte über `Authorization: Bearer …` oder `X-Brain-Key` geschützt. @@ -221,6 +227,7 @@ Wichtige Bereiche: - `ollama_pool.nodes[]`: Health, Kompatibilität, Inflight, Requests, Fehler und mittlere Laufzeit; - `glpi_kb`: letzter Sync, Dokumentanzahl, Pfad und Fehler; - `persistence`: ausstehende Dateien, Dirty-Status, letzter Flush und Flush-Fehler; +- `graph_storage`: SQLite-/WAL-Größe, Vector-Bytes, Row-Anzahl und ausstehende Änderungen; - `enrich_*`: Zustand des autonomen AI-THINK-Workers; - `relations_created`, `articles_created`, `articles_skipped`: getrennte Relation- und Artikelergebnisse. @@ -234,4 +241,14 @@ node --check internal/web/static/app.js ## Sichtbare SearXNG-Recherche -SearXNG-Suchen werden als eigenständige Aktivität visualisiert. Scanringe markieren die laufende Suche, gefundene Quellen erscheinen als externe Quellenpunkte und fließen anschließend in die neu erzeugten Forschungs-Nodes. Die Animation bleibt bei der Aktualisierung des Graphen bestehen und läuft nach dem Einblenden der neuen Nodes mindestens zwei Sekunden weiter. Der Aktivitätsfeed zeigt Trefferzahl, Domains, Quelltitel und Übernahmestatus. Details stehen in `SEARXNG-VISUALIZATION.md`. +SearXNG-Suchen werden als eigenständige Aktivität visualisiert. Scanringe markieren die laufende Suche, gefundene Quellen erscheinen als externe Quellenpunkte und fließen anschließend in die neu erzeugten Forschungs-Nodes. Die Animation bleibt bei der Aktualisierung des Graphen bestehen und läuft nach dem Einblenden der neuen Nodes mindestens zwei Sekunden weiter. + +Unter **FILTER → SearXNG-Diagnose** kann eine direkte Testsuche ausgeführt werden. Diese umgeht die fachliche Qwen-Entscheidung, verändert den Graphen nicht und zeigt eindeutig: + +- die tatsächlich verwendete Basis-URL, +- HTTP-Status und Content-Type, +- Trefferzahl und Laufzeit, +- DNS-, Netzwerk-, TLS-, HTTP- und JSON-Fehler, +- bei Fehlern den Antwortausschnitt von SearXNG oder Reverse Proxy. + +Die API-Endpunkte sind `GET /api/research/status` und `POST /api/research/test`. Der Aktivitätsfeed zeigt bei fehlgeschlagenen automatischen Suchen nun ebenfalls die konkrete Fehlerursache. Details stehen in `SEARXNG-VISUALIZATION.md`. diff --git a/SEARXNG-VISUALIZATION.md b/SEARXNG-VISUALIZATION.md index 754ff18..f6fd312 100644 --- a/SEARXNG-VISUALIZATION.md +++ b/SEARXNG-VISUALIZATION.md @@ -29,3 +29,72 @@ Der linke Feed zeigt jetzt: - Fehler und leere Ergebnismengen. Damit ist erkennbar, ob SearXNG tatsächlich Ergebnisse geliefert hat und ob diese in den Graphen übernommen wurden. + +## Direkter Funktionstest + +Im Einstellungsbereich unter **FILTER → SearXNG-Diagnose** steht eine direkte Testsuche zur Verfügung. Sie ruft SearXNG unabhängig von Qwens Entscheidung auf und übernimmt die Treffer nicht dauerhaft in den Graphen. Damit kann eindeutig zwischen zwei Fällen unterschieden werden: + +- SearXNG ist technisch nicht erreichbar oder liefert kein JSON. +- SearXNG funktioniert, aber Qwen hat für einen bestimmten Wissensverbund keine Recherche angefordert. + +Der Test erzeugt die Events: + +- `research.test.started` +- `research.test.results` +- `research.test.failed` + +Die normale Rechercheanimation wird dabei ebenfalls mindestens zwei Sekunden angezeigt. + +### Diagnoseinformationen + +`GET /api/research/status` liefert den zuletzt bekannten Zustand. `POST /api/research/test` akzeptiert beispielsweise: + +```json +{ + "query": "Btrfs Snapshots und ZFS History Unterschiede Timeline", + "limit": 4 +} +``` + +Die Antwort beziehungsweise das Fehler-Event enthält: + +```json +{ + "configured": true, + "base_url": "http://searxng:8080", + "ok": false, + "duration_ms": 29, + "http_status": 403, + "content_type": "text/html", + "error_kind": "http_status", + "error": "SearXNG lieferte HTTP 403 ..." +} +``` + +Im linken Aktivitätsfeed wird der vollständige Fehlertext einschließlich Endpoint angezeigt. Typische Fehlerklassen sind `dns`, `connection_refused`, `timeout`, `tls`, `http_status` und `invalid_json`. + +## Häufige Docker-Fehler + +`SEARXNG_URL=http://localhost:8080` ist in einem getrennten Brain-Container normalerweise falsch. `localhost` bezeichnet dort den Brain-Container selbst. Bei gemeinsamem Compose-Netzwerk sollte stattdessen beispielsweise gelten: + +```env +BRAIN_RESEARCH_ENABLED=true +SEARXNG_URL=http://searxng:8080 +``` + +Läuft SearXNG direkt auf dem Docker-Host, kann je nach Plattform verwendet werden: + +```env +SEARXNG_URL=http://host.docker.internal:8080 +``` + +SearXNG muss außerdem JSON-Antworten erlauben: + +```yaml +search: + formats: + - html + - json +``` + +Eine HTML-Seite, Login-Seite oder Proxy-Fehlerseite wird jetzt als `invalid_json` mit Content-Type und kurzem Antwortausschnitt gemeldet. diff --git a/SQLITE-STORAGE.md b/SQLITE-STORAGE.md new file mode 100644 index 0000000..4197bcf --- /dev/null +++ b/SQLITE-STORAGE.md @@ -0,0 +1,42 @@ +# Graph Store: modernc.org/sqlite + +## Ziel + +`graph-state.json` wurde vollständig durch SQLite ersetzt. Die Umstellung optimiert drei Engpässe großer Knowledge-Graphs: + +1. Embeddings werden nicht mehr als lange JSON-Dezimalzahlen gespeichert. +2. Ein Flush serialisiert nicht mehr den vollständigen Graphen. +3. Eine unveränderte KB-Synchronisation erzeugt keine erneuten Graphwrites. + +## Implementierung + +- Treiber: `modernc.org/sqlite v1.37.1` +- API: Go `database/sql` +- CGO: nicht erforderlich +- Datenbank: `BRAIN_DATA_DIR/graph.db` +- Vektoren: binäre `float32`-BLOBs +- Modellschutz: gespeicherter Name und Ollama-Digest; bei Änderung automatische Vector-Invalidierung +- Journal: WAL +- Flush: inkrementelle Upserts/Deletes innerhalb einer Transaktion +- Export: `VACUUM INTO` über `GET /api/state/export` + +`modernc.org/libc` ist passend zur verwendeten SQLite-Version auf `v1.65.7` festgesetzt. + +## Verhalten bei identischer KB-Basis + +Ein leistungsfähiges System kann den Graphen inklusive Embeddings und Beziehungen aufbauen. Der exportierte `graph.db` lässt sich anschließend auf ein schwächeres System kopieren. Wichtig sind stabile Dokument-IDs und dasselbe Embedding-Modell. + +Nach dem Import führt das Zielsystem weiterhin einen Quellenabgleich durch: + +- unverändert: Node, Edge und Vector bleiben bestehen; +- Metadatenänderung ohne Textänderung: nur die Node-Zeile wird aktualisiert; +- Titel, Summary, Kategorien oder Keywords geändert: der alte Vector wird gelöscht und selektiv neu erzeugt; +- Quelle entfernt: zugehörige Nodes, Vektoren und verwaiste Edges werden gelöscht. + +## Keine Legacy-Migration + +Die aktuelle Umgebung ist Staging. Deshalb gibt es absichtlich keinen Importpfad aus `graph-state.json`. Eine alte JSON-Datei wird im Status nur als `legacy_json_present` gemeldet und ansonsten ignoriert. + +## Speicherarme Startkonfiguration + +SQLite wird mit einem kleinen Page-Cache von ungefähr 8 MiB, deaktiviertem mmap und dateibasiertem temporärem Speicher geöffnet. Der vollständige Graph und die Embeddings liegen bereits im Go-Arbeitsspeicher; ein zusätzlicher großer SQLite-Cache würde den Speicherbedarf nur verdoppeln. Vor dem Öffnen prüft der Dienst außerdem, ob `BRAIN_DATA_DIR` tatsächlich beschreibbar ist. diff --git a/VALIDATION-SEARXNG-DIAGNOSTICS.md b/VALIDATION-SEARXNG-DIAGNOSTICS.md new file mode 100644 index 0000000..58f630f --- /dev/null +++ b/VALIDATION-SEARXNG-DIAGNOSTICS.md @@ -0,0 +1,25 @@ +# Validierung: SearXNG-Diagnose + +Ausgeführt: + +```text +node --check internal/web/static/app.js +go test ./internal/research +go test ./internal/engine -run TestResearchDiagnostic +go test ./internal/web -run TestResearchDiagnosticAPI +go test -run '^$' ./... +go vet ./... +CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build ./cmd/brain +Compose-YAML-Parsing +``` + +Die Research-Tests verwenden echte lokale HTTP-Testserver und prüfen: + +- erfolgreiche SearXNG-JSON-Antworten, +- Basis-URLs mit und ohne abschließendes `/search`, +- HTTP-Fehler einschließlich Antworttext, +- HTML statt JSON und den Hinweis auf `search.formats: json`, +- sichtbare Start-, Result- und Fehler-Events, +- die Diagnose-API. + +Da die Laufzeitumgebung keinen DNS-Zugriff auf `proxy.golang.org` hatte, wurden Pakete außerhalb des Research-Moduls für Typprüfung, Vet und Build mit einem lokalen compile-only Stub des bereits fest gepinnten `modernc.org/sqlite`-Imports geprüft. Der Stub ist nicht Bestandteil des Projekts oder ZIP-Archivs. Der Docker-Build verwendet weiterhin `modernc.org/sqlite v1.37.1` aus `go.mod` und `go.sum`. diff --git a/VALIDATION-SQLITE.md b/VALIDATION-SQLITE.md new file mode 100644 index 0000000..ae9be54 --- /dev/null +++ b/VALIDATION-SQLITE.md @@ -0,0 +1,16 @@ +# Validierung der SQLite-Umstellung + +Durchgeführt: + +- `gofmt` für alle Go-Quellen; +- vollständige Go-Typ- und Syntaxkompilierung aller Packages; +- `go vet ./...`; +- statischer `CGO_ENABLED=0`-Build der Anwendung; +- JavaScript-Syntaxprüfung des Webfrontends; +- YAML-Parsing beider Compose-Dateien; +- SQL-DDL-, Upsert-, Foreign-Key-, Float32-BLOB- und `VACUUM INTO`-Smoke-Test mit einer lokalen SQLite-Referenz; +- neue Tests für Roundtrip, inkrementelle Dirty Rows, portable Exporte, unveränderte KB-Scans, selektive Vector-Invalidierung und Embedding-Modellwechsel. + +Einschränkung dieser Build-Umgebung: + +Der Container konnte `proxy.golang.org` wegen gesperrter DNS-/Netzwerkauflösung nicht erreichen. Deshalb wurden die Tests hier für Typ- und Syntaxprüfung mit lokalen Modul-Stubs kompiliert; ein echter Testlauf gegen den `modernc.org/sqlite`-Treiber war in dieser Umgebung nicht möglich. Das Projekt referenziert die echten, fest gepinnten Module in `go.mod`. Beim Docker-Build lädt `go mod download` diese Abhängigkeiten und `go test ./...` führt die enthaltenen SQLite-Runtime-Tests aus. diff --git a/cmd/brain/main.go b/cmd/brain/main.go index ee7b0d2..7563097 100644 --- a/cmd/brain/main.go +++ b/cmd/brain/main.go @@ -40,7 +40,8 @@ func main() { ui := &webui.Server{Engine: eng, Graph: g, Broker: broker, APIKey: cfg.APIKey} srv := &http.Server{Addr: cfg.ListenAddr, Handler: ui.Handler(), ReadHeaderTimeout: 10 * time.Second, ReadTimeout: 30 * time.Second, WriteTimeout: 10 * time.Minute, IdleTimeout: 90 * time.Second} go func() { - slog.Info("neural brain listening", "addr", cfg.ListenAddr, "knowledge_dirs", cfg.KnowledgeDirs, "staging_dirs", cfg.StagingDirs) + storage := g.StorageStatus() + slog.Info("neural brain listening", "addr", cfg.ListenAddr, "knowledge_dirs", cfg.KnowledgeDirs, "staging_dirs", cfg.StagingDirs, "graph_backend", storage.Backend, "graph_db", storage.Path) if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { slog.Error("server failed", "error", err) cancel() @@ -49,6 +50,9 @@ func main() { <-ctx.Done() shutdown, c := context.WithTimeout(context.Background(), 10*time.Second) defer c() - _ = eng.Flush(shutdown) _ = srv.Shutdown(shutdown) + _ = eng.Flush(shutdown) + if err := g.Close(); err != nil { + slog.Error("graph close failed", "error", err) + } } diff --git a/data/.gitkeep b/data/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/data/graph-state.json b/data/graph-state.json deleted file mode 100644 index 319880c..0000000 --- a/data/graph-state.json +++ /dev/null @@ -1 +0,0 @@ -{"version":4,"nodes":null,"edges":null} \ No newline at end of file diff --git a/data/graph.db b/data/graph.db new file mode 100644 index 0000000..c4539b0 Binary files /dev/null and b/data/graph.db differ diff --git a/data/runtime-settings.json b/data/runtime-settings.json index 8f96550..679f015 100644 --- a/data/runtime-settings.json +++ b/data/runtime-settings.json @@ -1,8 +1,9 @@ { "learning_enabled": false, - "thinking_enabled": false, + "thinking_enabled": true, "learning_categories": [], "display_categories": [], "thinking_categories": [], - "view_mode": "neural" + "view_mode": "neural", + "max_display_nodes": 1000 } diff --git a/deployment/docker-compose.full.yml b/deployment/docker-compose.full.yml index a63d24e..16e0f14 100644 --- a/deployment/docker-compose.full.yml +++ b/deployment/docker-compose.full.yml @@ -156,6 +156,8 @@ services: BRAIN_TOP_K: ${BRAIN_TOP_K:-8} BRAIN_MAX_CONTEXT_CHARS: ${BRAIN_MAX_CONTEXT_CHARS:-16000} BRAIN_RESEARCH_ENABLED: ${BRAIN_RESEARCH_ENABLED:-false} + # In Docker, localhost points to this Brain container. Use the SearXNG + # service name (for example http://searxng:8080) or host.docker.internal. SEARXNG_URL: ${SEARXNG_URL:-} GLPI_KB_ENABLED: ${GLPI_KB_ENABLED:-false} GLPI_URL: ${GLPI_URL:-} @@ -176,6 +178,8 @@ services: - ${KNOWLEDGE_HOST_PATH:-../../glpi-ai-agent/knowledge}:/sources/knowledge:ro - ${KB_STAGING_HOST_PATH:-../../glpi-ai-knowledgebase/staging}:/sources/staging:rw - agent-data:/sources/agent-data:ro + extra_hosts: + - "host.docker.internal:host-gateway" ollama: image: ollama/ollama:latest diff --git a/docker-compose.yml b/docker-compose.yml index 4526c69..a5375ab 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -55,6 +55,8 @@ services: BRAIN_TOP_K: ${BRAIN_TOP_K:-8} BRAIN_MAX_CONTEXT_CHARS: ${BRAIN_MAX_CONTEXT_CHARS:-16000} BRAIN_RESEARCH_ENABLED: ${BRAIN_RESEARCH_ENABLED:-false} + # In Docker, localhost points to this Brain container. Use the SearXNG + # service name (for example http://searxng:8080) or host.docker.internal. SEARXNG_URL: ${SEARXNG_URL:-} BRAIN_API_KEY: ${BRAIN_API_KEY:-} GLPI_KB_ENABLED: ${GLPI_KB_ENABLED:-false} diff --git a/go.mod b/go.mod index 163c11b..3841f3d 100644 --- a/go.mod +++ b/go.mod @@ -1,3 +1,18 @@ module github.com/local/glpi-neural-brain go 1.26 + +require modernc.org/sqlite v1.37.1 + +require ( + github.com/dustin/go-humanize v1.0.1 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + github.com/ncruces/go-strftime v0.1.9 // indirect + github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect + golang.org/x/exp v0.0.0-20250408133849-7e4ce0ab07d0 // indirect + golang.org/x/sys v0.33.0 // indirect + modernc.org/libc v1.65.7 // indirect + modernc.org/mathutil v1.7.1 // indirect + modernc.org/memory v1.11.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..c0ce562 --- /dev/null +++ b/go.sum @@ -0,0 +1,47 @@ +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e h1:ijClszYn+mADRFY17kjQEVQ1XRhq2/JR1M3sGqeJoxs= +github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e/go.mod h1:boTsfXsheKC2y+lKOCMpSfarhxDeIzfZG1jqGcPl3cA= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/ncruces/go-strftime v0.1.9 h1:bY0MQC28UADQmHmaF5dgpLmImcShSi2kHU9XLdhx/f4= +github.com/ncruces/go-strftime v0.1.9/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= +golang.org/x/exp v0.0.0-20250408133849-7e4ce0ab07d0 h1:R84qjqJb5nVJMxqWYb3np9L5ZsaDtB+a39EqjV0JSUM= +golang.org/x/exp v0.0.0-20250408133849-7e4ce0ab07d0/go.mod h1:S9Xr4PYopiDyqSyp5NjCrhFrqg6A5zA2E/iPHPhqnS8= +golang.org/x/mod v0.24.0 h1:ZfthKaKaT4NrhGVZHO1/WDTwGES4De8KtWO0SIbNJMU= +golang.org/x/mod v0.24.0/go.mod h1:IXM97Txy2VM4PJ3gI61r1YEk/gAj6zAHN3AdZt6S9Ww= +golang.org/x/sync v0.14.0 h1:woo0S4Yywslg6hp4eUFjTVOyKt0RookbpAHG4c1HmhQ= +golang.org/x/sync v0.14.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= +golang.org/x/sys v0.33.0 h1:q3i8TbbEz+JRD9ywIRlyRAQbM0qF7hu24q3teo2hbuw= +golang.org/x/sys v0.33.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/tools v0.33.0 h1:4qz2S3zmRxbGIhDIAgjxvFutSvH5EfnsYrRBj0UI0bc= +golang.org/x/tools v0.33.0/go.mod h1:CIJMaWEY88juyUfo7UbgPqbC8rU2OqfAV1h2Qp0oMYI= +modernc.org/cc/v4 v4.26.1 h1:+X5NtzVBn0KgsBCBe+xkDC7twLb/jNVj9FPgiwSQO3s= +modernc.org/cc/v4 v4.26.1/go.mod h1:uVtb5OGqUKpoLWhqwNQo/8LwvoiEBLvZXIQ/SmO6mL0= +modernc.org/ccgo/v4 v4.28.0 h1:rjznn6WWehKq7dG4JtLRKxb52Ecv8OUGah8+Z/SfpNU= +modernc.org/ccgo/v4 v4.28.0/go.mod h1:JygV3+9AV6SmPhDasu4JgquwU81XAKLd3OKTUDNOiKE= +modernc.org/fileutil v1.3.1 h1:8vq5fe7jdtEvoCf3Zf9Nm0Q05sH6kGx0Op2CPx1wTC8= +modernc.org/fileutil v1.3.1/go.mod h1:HxmghZSZVAz/LXcMNwZPA/DRrQZEVP9VX0V4LQGQFOc= +modernc.org/gc/v2 v2.6.5 h1:nyqdV8q46KvTpZlsw66kWqwXRHdjIlJOhG6kxiV/9xI= +modernc.org/gc/v2 v2.6.5/go.mod h1:YgIahr1ypgfe7chRuJi2gD7DBQiKSLMPgBQe9oIiito= +modernc.org/libc v1.65.7 h1:Ia9Z4yzZtWNtUIuiPuQ7Qf7kxYrxP1/jeHZzG8bFu00= +modernc.org/libc v1.65.7/go.mod h1:011EQibzzio/VX3ygj1qGFt5kMjP0lHb0qCW5/D/pQU= +modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU= +modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg= +modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI= +modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw= +modernc.org/opt v0.1.4 h1:2kNGMRiUjrp4LcaPuLY2PzUfqM/w9N23quVwhKt5Qm8= +modernc.org/opt v0.1.4/go.mod h1:03fq9lsNfvkYSfxrfUhZCWPk1lm4cq4N+Bh//bEtgns= +modernc.org/sortutil v1.2.1 h1:+xyoGf15mM3NMlPDnFqrteY07klSFxLElE2PVuWIJ7w= +modernc.org/sortutil v1.2.1/go.mod h1:7ZI3a3REbai7gzCLcotuw9AC4VZVpYMjDzETGsSMqJE= +modernc.org/sqlite v1.37.1 h1:EgHJK/FPoqC+q2YBXg7fUmES37pCHFc97sI7zSayBEs= +modernc.org/sqlite v1.37.1/go.mod h1:XwdRtsE1MpiBcL54+MbKcaDvcuej+IYSMfLN6gSKV8g= +modernc.org/strutil v1.2.1 h1:UneZBkQA+DX2Rp35KcM69cSsNES9ly8mQWD71HKlOA0= +modernc.org/strutil v1.2.1/go.mod h1:EHkiggD70koQxjVdSBM3JKM7k6L0FbGE5eymy9i3B9A= +modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y= +modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM= diff --git a/internal/engine/article.go b/internal/engine/article.go index 9ff41a5..ec2bbec 100644 --- a/internal/engine/article.go +++ b/internal/engine/article.go @@ -6,6 +6,7 @@ import ( "encoding/hex" "encoding/json" "fmt" + "log/slog" "math" "os" "path/filepath" @@ -467,12 +468,16 @@ func (e *Engine) researchKnowledgeGaps(ctx context.Context, trigger string, node if resultLimit < 1 { resultLimit = 4 } - results, err := e.Research.Search(ctx, query, resultLimit) + results, diagnostic, err := e.Research.SearchDetailed(ctx, query, resultLimit) if err != nil { - e.Broker.Publish(model.Activity{Type: "article.research.failed", Source: "searxng", Phase: "knowledge-research", NodeIDs: nodeIDs, Message: "Die ergänzende Artikelrecherche ist fehlgeschlagen", Strength: .35, Metadata: mergeResearchMetadata(startMetadata, map[string]any{"error": err.Error(), "duration_ms": time.Since(started).Milliseconds()})}) + metadata := mergeResearchMetadata(startMetadata, researchDiagnosticMetadata(diagnostic)) + metadata["error"] = err.Error() + metadata["duration_ms"] = time.Since(started).Milliseconds() + slog.Warn("article research failed", "query", query, "base_url", diagnostic.BaseURL, "kind", diagnostic.ErrorKind, "http_status", diagnostic.HTTPStatus, "duration_ms", diagnostic.DurationMS, "error", err) + e.Broker.Publish(model.Activity{Type: "article.research.failed", Source: "searxng", Phase: "knowledge-research", NodeIDs: nodeIDs, Message: "Die ergänzende Artikelrecherche ist fehlgeschlagen", Strength: .35, Metadata: metadata}) return nil, err } - resultMetadata := researchEventMetadata(trigger, researchID, query, results, time.Since(started)) + resultMetadata := mergeResearchMetadata(researchEventMetadata(trigger, researchID, query, results, time.Since(started)), researchDiagnosticMetadata(diagnostic)) message := fmt.Sprintf("SearXNG hat %d Quellen für den offenen Wissenspunkt geliefert", len(results)) if len(results) == 0 { message = "SearXNG hat für den offenen Wissenspunkt keine verwertbare Quelle geliefert" diff --git a/internal/engine/engine.go b/internal/engine/engine.go index f4698b9..f623346 100644 --- a/internal/engine/engine.go +++ b/internal/engine/engine.go @@ -124,6 +124,10 @@ func New(cfg config.Config, g *graph.Store, b *activity.Broker) *Engine { } nodes = append(nodes, ollama.NodeConfig{Name: name, URL: rawURL, Weight: weight}) } + clearedVectors := g.ConfigureEmbeddingModel(cfg.EmbeddingModel) + if clearedVectors > 0 && b != nil { + b.Publish(model.Activity{Type: "embedding.model_changed", Source: "brain", Phase: "learning", Message: fmt.Sprintf("Embedding-Modell geändert · %d Vektoren werden neu gelernt", clearedVectors), Strength: .65, Metadata: map[string]any{"embedding_model": cfg.EmbeddingModel, "cleared_vectors": clearedVectors}}) + } pool := ollama.NewPool(ollama.PoolConfig{ Nodes: nodes, RoutingMode: cfg.OllamaRoutingMode, NodeMaxInflight: cfg.OllamaNodeMaxInflight, HealthInterval: cfg.OllamaHealthInterval, FailureCooldown: cfg.OllamaFailureCooldown, @@ -361,16 +365,23 @@ func (e *Engine) idle(ctx context.Context) { case <-ctx.Done(): return case <-ticker.C: - s := e.Graph.Snapshot() - if len(s.Nodes) == 0 { + n, ok := e.Graph.IdleNode(time.Now().Unix() / 7) + if !ok { continue } - idx := int(time.Now().Unix()/7) % len(s.Nodes) - n := s.Nodes[idx] e.Broker.Publish(model.Activity{Type: "brain.idle", Source: "brain", Phase: "idle", Message: "Leise Hintergrundaktivität", NodeIDs: []string{n.ID}, Strength: .18}) } } } +func (e *Engine) configureEmbeddingDigest() int { + for _, node := range e.Ollama.NodeStatuses() { + if node.Healthy && node.Compatible && node.EmbeddingModel && strings.TrimSpace(node.EmbeddingDigest) != "" { + return e.Graph.ConfigureEmbeddingIdentity(e.Cfg.EmbeddingModel, node.EmbeddingDigest) + } + } + return 0 +} + func (e *Engine) Scan(ctx context.Context) error { if !e.LearningEnabled() { return ErrLearningDisabled @@ -394,6 +405,9 @@ func (e *Engine) Scan(ctx context.Context) error { e.ensureFallbackEmbeddings() e.setOllamaOK(false) } else { + if cleared := e.configureEmbeddingDigest(); cleared > 0 { + e.Broker.Publish(model.Activity{Type: "embedding.identity_changed", Source: "brain", Phase: "learning", Message: fmt.Sprintf("Embedding-Digest geändert · %d Vektoren werden neu gelernt", cleared), Strength: .7, Metadata: map[string]any{"embedding_model": e.Cfg.EmbeddingModel, "cleared_vectors": cleared}}) + } // Local fallback vectors use 256 dimensions. Once Ollama becomes available, // discard those placeholders and replace them with real model embeddings. e.Graph.ClearVectorsByDimension(256) @@ -409,8 +423,8 @@ func (e *Engine) Scan(ctx context.Context) error { e.lastScan = time.Now().UTC() e.stateMu.Unlock() if e.Graph.Version() != beforeVersion { - s := e.Graph.Snapshot() - e.Broker.Publish(model.Activity{Type: "graph.updated", Source: "brain", Phase: "indexed", Message: fmt.Sprintf("%d Wissenselemente · %d Nodes · %d Edges", count, len(s.Nodes), len(s.Edges)), Strength: .55, Metadata: map[string]any{"nodes": len(s.Nodes), "edges": len(s.Edges), "knowledge_elements": count}}) + nodes, edges, _ := e.Graph.Counts() + e.Broker.Publish(model.Activity{Type: "graph.updated", Source: "brain", Phase: "indexed", Message: fmt.Sprintf("%d Wissenselemente · %d Nodes · %d Edges", count, nodes, edges), Strength: .55, Metadata: map[string]any{"nodes": nodes, "edges": edges, "knowledge_elements": count}}) } return nil } @@ -631,12 +645,15 @@ func (e *Engine) enrichOne(ctx context.Context, trigger string) (EnrichOutcome, researchStarted := time.Now() startMetadata := map[string]any{"trigger": trigger, "research_id": researchID, "research_query": decision.ResearchQuery, "source_label": a.Label, "target_label": b.Label, "animation_min_ms": 2000} e.Broker.Publish(model.Activity{Type: "research.started", Source: "searxng", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: "Unklarheit erkannt · SearXNG durchsucht externe Quellen", Strength: .9, Metadata: startMetadata}) - results, err := e.Research.Search(ctx, decision.ResearchQuery, 4) + results, diagnostic, err := e.Research.SearchDetailed(ctx, decision.ResearchQuery, 4) if err != nil { - slog.Warn("research failed", "error", err) - e.Broker.Publish(model.Activity{Type: "research.failed", Source: "searxng", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: "SearXNG-Recherche ist fehlgeschlagen", Strength: .35, Metadata: mergeResearchMetadata(startMetadata, map[string]any{"error": err.Error(), "duration_ms": time.Since(researchStarted).Milliseconds()})}) + metadata := mergeResearchMetadata(startMetadata, researchDiagnosticMetadata(diagnostic)) + metadata["error"] = err.Error() + metadata["duration_ms"] = time.Since(researchStarted).Milliseconds() + slog.Warn("research failed", "query", decision.ResearchQuery, "base_url", diagnostic.BaseURL, "kind", diagnostic.ErrorKind, "http_status", diagnostic.HTTPStatus, "duration_ms", diagnostic.DurationMS, "error", err) + e.Broker.Publish(model.Activity{Type: "research.failed", Source: "searxng", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: "SearXNG-Recherche ist fehlgeschlagen", Strength: .35, Metadata: metadata}) } else { - resultMetadata := researchEventMetadata(trigger, researchID, decision.ResearchQuery, results, time.Since(researchStarted)) + resultMetadata := mergeResearchMetadata(researchEventMetadata(trigger, researchID, decision.ResearchQuery, results, time.Since(researchStarted)), researchDiagnosticMetadata(diagnostic)) message := fmt.Sprintf("SearXNG hat %d verwertbare Webquellen geliefert", len(results)) if len(results) == 0 { message = "SearXNG hat keine verwertbaren Webquellen geliefert" @@ -706,10 +723,10 @@ func (e *Engine) addResearch(a, b model.Node, results []model.ResearchResult) re return uniqueResearchRefs(refs) } func (e *Engine) Status() map[string]any { - s := e.Graph.Snapshot() + nodes, edges, version := e.Graph.Counts() e.stateMu.RLock() status := map[string]any{ - "ok": true, "nodes": len(s.Nodes), "edges": len(s.Edges), "version": s.Version, + "ok": true, "nodes": nodes, "edges": edges, "version": version, "last_scan": e.lastScan, "last_enrich": e.lastEnrich, "last_enrich_attempt": e.lastAttempt, "next_enrich": e.nextEnrich, "ollama_ok": e.ollamaOK, "auto_enrich": e.Cfg.AutoEnrich, "enrich_running": e.enrichRunning, "enrich_trigger": e.enrichTrigger, "enrich_result": e.enrichResult, @@ -722,7 +739,8 @@ func (e *Engine) Status() map[string]any { "enrich_interval": e.Cfg.EnrichInterval.String(), "enrich_batch_size": e.Cfg.EnrichBatchSize, "enrich_anchors": e.Cfg.EnrichAnchors, "research_enabled": e.Cfg.ResearchEnabled, "chat_model": e.Cfg.ChatModel, "embedding_model": e.Cfg.EmbeddingModel, - "ollama_pool": e.Ollama.PoolStatus(), "persistence": e.Persistence.Status(), + "searxng": e.ResearchStatus(), + "ollama_pool": e.Ollama.PoolStatus(), "persistence": e.Persistence.Status(), "graph_storage": e.Graph.StorageStatus(), "runtime_settings": e.RuntimeSettings(), } if e.GLPIKB != nil { @@ -738,6 +756,13 @@ func (e *Engine) Flush(ctx context.Context) error { return e.Persistence.Flush(ctx, "manual") } +func (e *Engine) ExportGraph(ctx context.Context, destination string) error { + if err := e.Persistence.Flush(ctx, "export"); err != nil { + return err + } + return e.Graph.Export(ctx, destination) +} + func (e *Engine) SyncGLPIKB(ctx context.Context) error { if !e.LearningEnabled() { return ErrLearningDisabled diff --git a/internal/engine/research_diagnostics.go b/internal/engine/research_diagnostics.go new file mode 100644 index 0000000..3e5390b --- /dev/null +++ b/internal/engine/research_diagnostics.go @@ -0,0 +1,103 @@ +package engine + +import ( + "context" + "errors" + "fmt" + "log/slog" + "strings" + "time" + + "github.com/local/glpi-neural-brain/internal/model" + "github.com/local/glpi-neural-brain/internal/research" +) + +type ResearchTestResult struct { + OK bool `json:"ok"` + Query string `json:"query"` + Diagnostic research.Diagnostic `json:"diagnostic"` + Results []model.ResearchResult `json:"results"` +} + +func (e *Engine) ResearchStatus() research.Diagnostic { + if e == nil || e.Research == nil { + return research.Diagnostic{Configured: false, OK: false, ErrorKind: "not_configured", Error: "SearXNG ist nicht konfiguriert"} + } + status := e.Research.Status() + if !e.Cfg.ResearchEnabled { + status.OK = false + status.ErrorKind = "disabled" + status.Error = "BRAIN_RESEARCH_ENABLED ist deaktiviert" + } + return status +} + +func (e *Engine) TestResearch(ctx context.Context, query string, limit int) (ResearchTestResult, error) { + query = strings.TrimSpace(query) + if query == "" { + query = "Btrfs Snapshots und ZFS History Unterschiede Timeline" + } + if limit < 1 { + limit = 4 + } + if limit > 10 { + limit = 10 + } + if !e.Cfg.ResearchEnabled { + err := errors.New("SearXNG-Recherche ist deaktiviert: BRAIN_RESEARCH_ENABLED=false") + return ResearchTestResult{Query: query, Diagnostic: e.ResearchStatus()}, err + } + if e.Research == nil { + err := errors.New("SearXNG-Client ist nicht konfiguriert") + return ResearchTestResult{Query: query, Diagnostic: e.ResearchStatus()}, err + } + + researchID := newResearchRunID("research-test", query) + started := time.Now() + startMetadata := map[string]any{ + "trigger": "diagnostic", + "research_id": researchID, + "research_query": query, + "animation_min_ms": 2000, + "test_only": true, + "searxng_base_url": e.Research.Status().BaseURL, + } + e.Broker.Publish(model.Activity{Type: "research.test.started", Source: "searxng", Phase: "diagnostic", Message: "SearXNG-Verbindung und JSON-Suche werden direkt getestet", Strength: .92, Metadata: startMetadata}) + + results, diagnostic, err := e.Research.SearchDetailed(ctx, query, limit) + if err != nil { + metadata := mergeResearchMetadata(startMetadata, researchDiagnosticMetadata(diagnostic)) + metadata["duration_ms"] = time.Since(started).Milliseconds() + metadata["error"] = err.Error() + slog.Warn("SearXNG diagnostic search failed", "base_url", diagnostic.BaseURL, "kind", diagnostic.ErrorKind, "duration_ms", diagnostic.DurationMS, "error", err) + e.Broker.Publish(model.Activity{Type: "research.test.failed", Source: "searxng", Phase: "diagnostic", Message: "SearXNG-Test fehlgeschlagen", Strength: .4, Metadata: metadata}) + return ResearchTestResult{OK: false, Query: query, Diagnostic: diagnostic}, err + } + + metadata := mergeResearchMetadata(researchEventMetadata("diagnostic", researchID, query, results, time.Since(started)), researchDiagnosticMetadata(diagnostic)) + metadata["test_only"] = true + message := fmt.Sprintf("SearXNG-Test erfolgreich · %d verwertbare Treffer", len(results)) + if len(results) == 0 { + message = "SearXNG-Test erfolgreich · JSON-Antwort erhalten, aber keine verwertbaren Treffer" + } + e.Broker.Publish(model.Activity{Type: "research.test.results", Source: "searxng", Phase: "diagnostic", Message: message, Strength: 1, Metadata: metadata}) + return ResearchTestResult{OK: true, Query: query, Diagnostic: diagnostic, Results: results}, nil +} + +func researchDiagnosticMetadata(diagnostic research.Diagnostic) map[string]any { + metadata := map[string]any{ + "searxng_configured": diagnostic.Configured, + "searxng_ok": diagnostic.OK, + "searxng_base_url": diagnostic.BaseURL, + "searxng_http_status": diagnostic.HTTPStatus, + "searxng_content_type": diagnostic.ContentType, + "error_kind": diagnostic.ErrorKind, + } + if diagnostic.DurationMS > 0 { + metadata["duration_ms"] = diagnostic.DurationMS + } + if diagnostic.Error != "" { + metadata["error"] = diagnostic.Error + } + return metadata +} diff --git a/internal/engine/research_diagnostics_test.go b/internal/engine/research_diagnostics_test.go new file mode 100644 index 0000000..8e568a9 --- /dev/null +++ b/internal/engine/research_diagnostics_test.go @@ -0,0 +1,70 @@ +package engine + +import ( + "context" + "net/http" + "net/http/httptest" + "testing" + + "github.com/local/glpi-neural-brain/internal/activity" + "github.com/local/glpi-neural-brain/internal/config" + "github.com/local/glpi-neural-brain/internal/research" +) + +func TestResearchDiagnosticPublishesVisibleEvents(t *testing.T) { + searx := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"results":[{"title":"ZFS history","url":"https://example.test/zfs","content":"Timeline evidence"}]}`)) + })) + defer searx.Close() + + broker := activity.New(20) + eng := &Engine{ + Cfg: config.Config{ResearchEnabled: true}, + Research: research.New(searx.URL), + Broker: broker, + } + result, err := eng.TestResearch(context.Background(), "zfs timeline", 4) + if err != nil { + t.Fatal(err) + } + if !result.OK || len(result.Results) != 1 || !result.Diagnostic.OK { + t.Fatalf("unexpected result: %+v", result) + } + events := broker.Recent() + if len(events) != 2 || events[0].Type != "research.test.started" || events[1].Type != "research.test.results" { + t.Fatalf("unexpected events: %+v", events) + } + if events[1].Metadata["result_count"] != 1 { + t.Fatalf("result count missing from event: %+v", events[1].Metadata) + } +} + +func TestResearchDiagnosticIncludesFailureReason(t *testing.T) { + searx := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusForbidden) + _, _ = w.Write([]byte("json disabled")) + })) + defer searx.Close() + + broker := activity.New(20) + eng := &Engine{ + Cfg: config.Config{ResearchEnabled: true}, + Research: research.New(searx.URL), + Broker: broker, + } + result, err := eng.TestResearch(context.Background(), "zfs timeline", 4) + if err == nil { + t.Fatal("expected failure") + } + if result.Diagnostic.ErrorKind != "http_status" || result.Diagnostic.HTTPStatus != http.StatusForbidden { + t.Fatalf("unexpected diagnostic: %+v", result.Diagnostic) + } + events := broker.Recent() + if len(events) != 2 || events[1].Type != "research.test.failed" { + t.Fatalf("unexpected events: %+v", events) + } + if events[1].Metadata["error_kind"] != "http_status" || events[1].Metadata["searxng_http_status"] != http.StatusForbidden { + t.Fatalf("failure details missing: %+v", events[1].Metadata) + } +} diff --git a/internal/graph/sqlite_backend.go b/internal/graph/sqlite_backend.go new file mode 100644 index 0000000..2ddb867 --- /dev/null +++ b/internal/graph/sqlite_backend.go @@ -0,0 +1,700 @@ +package graph + +import ( + "context" + "database/sql" + "encoding/binary" + "encoding/json" + "errors" + "fmt" + "math" + "os" + "path/filepath" + "strings" + "time" + + "github.com/local/glpi-neural-brain/internal/model" + _ "modernc.org/sqlite" +) + +const schemaVersion = 1 + +const nodeUpsertSQL = `INSERT INTO nodes(id,kind,label,summary,status,origin,external_id,uri,categories_json,keywords_json,metadata_json,weight,x,y,z,updated_at_ns) +VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) +ON CONFLICT(id) DO UPDATE SET kind=excluded.kind,label=excluded.label,summary=excluded.summary,status=excluded.status,origin=excluded.origin,external_id=excluded.external_id,uri=excluded.uri,categories_json=excluded.categories_json,keywords_json=excluded.keywords_json,metadata_json=excluded.metadata_json,weight=excluded.weight,x=excluded.x,y=excluded.y,z=excluded.z,updated_at_ns=excluded.updated_at_ns` + +const edgeUpsertSQL = `INSERT INTO edges(id,source,target,type,origin,status,confidence,weight,explanation,evidence_json,metadata_json,created_at_ns,updated_at_ns) +VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?) +ON CONFLICT(id) DO UPDATE SET source=excluded.source,target=excluded.target,type=excluded.type,origin=excluded.origin,status=excluded.status,confidence=excluded.confidence,weight=excluded.weight,explanation=excluded.explanation,evidence_json=excluded.evidence_json,metadata_json=excluded.metadata_json,created_at_ns=excluded.created_at_ns,updated_at_ns=excluded.updated_at_ns` + +const vectorUpsertSQL = `INSERT INTO vectors(node_id,dimensions,data,updated_at_ns) VALUES(?,?,?,?) +ON CONFLICT(node_id) DO UPDATE SET dimensions=excluded.dimensions,data=excluded.data,updated_at_ns=excluded.updated_at_ns` + +// StorageStatus describes the persistent SQLite graph store. The live graph is +// still held in memory; only changed rows are written during a flush. +type StorageStatus struct { + Backend string `json:"backend"` + Path string `json:"path"` + SchemaVersion int `json:"schema_version"` + JournalMode string `json:"journal_mode"` + DatabaseBytes int64 `json:"database_bytes"` + WALBytes int64 `json:"wal_bytes"` + NodeRows int `json:"node_rows"` + EdgeRows int `json:"edge_rows"` + VectorRows int `json:"vector_rows"` + VectorBytes int64 `json:"vector_bytes"` + PersistedVersion uint64 `json:"persisted_version"` + EmbeddingModel string `json:"embedding_model,omitempty"` + EmbeddingDigest string `json:"embedding_digest,omitempty"` + LegacyJSONPresent bool `json:"legacy_json_present"` + PendingNodes int `json:"pending_nodes"` + PendingEdges int `json:"pending_edges"` + PendingVectors int `json:"pending_vectors"` + PendingDeletions int `json:"pending_deletions"` + PendingMetadata bool `json:"pending_metadata"` +} + +func Open(dir string) (*Store, error) { + absoluteDir, err := filepath.Abs(dir) + if err != nil { + return nil, fmt.Errorf("resolve graph data directory: %w", err) + } + if err := os.MkdirAll(absoluteDir, 0o750); err != nil { + return nil, fmt.Errorf("create graph data directory %q: %w", absoluteDir, err) + } + if err := ensureWritableDirectory(absoluteDir); err != nil { + return nil, err + } + + path := filepath.Join(absoluteDir, "graph.db") + if info, statErr := os.Stat(path); statErr == nil && info.IsDir() { + return nil, fmt.Errorf("graph database path %q is a directory; remove or rename it", path) + } else if statErr != nil && !errors.Is(statErr, os.ErrNotExist) { + return nil, fmt.Errorf("inspect graph database %q: %w", path, statErr) + } + + // Use the plain filename here instead of applying PRAGMAs through the DSN. + // This makes startup failures attributable to one concrete step and avoids + // asking the driver to reserve a large mmap/cache while the connection is + // being created. The application keeps the full graph in Go memory already, + // so an additional 64 MiB SQLite page cache and 256 MiB mmap are wasteful on + // small systems. + db, err := sql.Open("sqlite", path) + if err != nil { + return nil, fmt.Errorf("create sqlite handle for %q: %w", path, err) + } + // The Store serializes persistence and reads the live graph from memory. + // One SQLite connection is sufficient and keeps all connection-local + // PRAGMAs deterministic. + db.SetMaxOpenConns(1) + db.SetMaxIdleConns(1) + db.SetConnMaxLifetime(0) + db.SetConnMaxIdleTime(0) + + openCtx, openCancel := context.WithTimeout(context.Background(), 15*time.Second) + if err := db.PingContext(openCtx); err != nil { + openCancel() + db.Close() + return nil, fmt.Errorf("open sqlite database %q: %w", path, err) + } + journalMode, err := configureSQLite(openCtx, db) + openCancel() + if err != nil { + db.Close() + return nil, fmt.Errorf("configure sqlite database %q: %w", path, err) + } + + s := &Store{ + nodes: map[string]model.Node{}, + edges: map[string]model.Edge{}, + vectors: map[string][]float32{}, + db: db, + dbPath: path, + journalMode: journalMode, + dirtyNodes: map[string]uint64{}, + dirtyEdges: map[string]uint64{}, + dirtyVectors: map[string]uint64{}, + deletedNodes: map[string]uint64{}, + deletedEdges: map[string]uint64{}, + deletedVectors: map[string]uint64{}, + } + loadCtx, loadCancel := context.WithTimeout(context.Background(), 10*time.Minute) + defer loadCancel() + if err := s.initSchema(loadCtx); err != nil { + db.Close() + return nil, fmt.Errorf("initialize sqlite graph schema: %w", err) + } + if err := s.load(loadCtx); err != nil { + db.Close() + return nil, fmt.Errorf("load sqlite graph state: %w", err) + } + return s, nil +} + +func ensureWritableDirectory(dir string) error { + info, err := os.Stat(dir) + if err != nil { + return fmt.Errorf("inspect graph data directory %q: %w", dir, err) + } + if !info.IsDir() { + return fmt.Errorf("graph data path %q is not a directory", dir) + } + probe, err := os.CreateTemp(dir, ".brain-write-test-*") + if err != nil { + return fmt.Errorf("graph data directory %q is not writable by the brain process: %w", dir, err) + } + name := probe.Name() + closeErr := probe.Close() + removeErr := os.Remove(name) + if closeErr != nil { + return fmt.Errorf("close graph data write test %q: %w", name, closeErr) + } + if removeErr != nil { + return fmt.Errorf("remove graph data write test %q: %w", name, removeErr) + } + return nil +} + +func configureSQLite(ctx context.Context, db *sql.DB) (string, error) { + // Conservative defaults: the live graph and vectors already reside in Go + // memory. SQLite therefore only needs a small page cache and no mmap window. + // This keeps startup predictable inside memory-limited containers. + statements := []struct { + name string + sql string + }{ + {name: "busy_timeout", sql: `PRAGMA busy_timeout=5000`}, + {name: "foreign_keys", sql: `PRAGMA foreign_keys=ON`}, + {name: "synchronous", sql: `PRAGMA synchronous=NORMAL`}, + {name: "temp_store", sql: `PRAGMA temp_store=FILE`}, + {name: "cache_size", sql: `PRAGMA cache_size=-8192`}, + {name: "mmap_size", sql: `PRAGMA mmap_size=0`}, + {name: "wal_autocheckpoint", sql: `PRAGMA wal_autocheckpoint=0`}, + {name: "journal_size_limit", sql: `PRAGMA journal_size_limit=67108864`}, + } + for _, statement := range statements { + if _, err := db.ExecContext(ctx, statement.sql); err != nil { + return "", fmt.Errorf("apply PRAGMA %s: %w", statement.name, err) + } + } + + var journalMode string + if err := db.QueryRowContext(ctx, `PRAGMA journal_mode=WAL`).Scan(&journalMode); err != nil { + return "", fmt.Errorf("enable WAL journal mode: %w", err) + } + journalMode = strings.ToLower(strings.TrimSpace(journalMode)) + if journalMode != "wal" { + return "", fmt.Errorf("enable WAL journal mode: SQLite selected %q", journalMode) + } + return journalMode, nil +} + +func (s *Store) initSchema(ctx context.Context) error { + statements := []string{ + `CREATE TABLE IF NOT EXISTS graph_meta ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL +) WITHOUT ROWID`, + `CREATE TABLE IF NOT EXISTS nodes ( + id TEXT PRIMARY KEY, + kind TEXT NOT NULL, + label TEXT NOT NULL, + summary TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT '', + origin TEXT NOT NULL, + external_id TEXT NOT NULL DEFAULT '', + uri TEXT NOT NULL DEFAULT '', + categories_json TEXT NOT NULL DEFAULT '[]', + keywords_json TEXT NOT NULL DEFAULT '[]', + metadata_json TEXT NOT NULL DEFAULT '{}', + weight REAL NOT NULL DEFAULT 1, + x REAL NOT NULL DEFAULT 0, + y REAL NOT NULL DEFAULT 0, + z REAL NOT NULL DEFAULT 0, + updated_at_ns INTEGER NOT NULL +) WITHOUT ROWID`, + `CREATE INDEX IF NOT EXISTS idx_nodes_origin ON nodes(origin)`, + `CREATE INDEX IF NOT EXISTS idx_nodes_kind_status ON nodes(kind, status)`, + `CREATE INDEX IF NOT EXISTS idx_nodes_external_id ON nodes(external_id)`, + `CREATE TABLE IF NOT EXISTS edges ( + id TEXT PRIMARY KEY, + source TEXT NOT NULL REFERENCES nodes(id) ON DELETE CASCADE, + target TEXT NOT NULL REFERENCES nodes(id) ON DELETE CASCADE, + type TEXT NOT NULL, + origin TEXT NOT NULL, + status TEXT NOT NULL DEFAULT '', + confidence REAL NOT NULL DEFAULT 0, + weight REAL NOT NULL DEFAULT 1, + explanation TEXT NOT NULL DEFAULT '', + evidence_json TEXT NOT NULL DEFAULT '[]', + metadata_json TEXT NOT NULL DEFAULT '{}', + created_at_ns INTEGER NOT NULL, + updated_at_ns INTEGER NOT NULL +) WITHOUT ROWID`, + `CREATE INDEX IF NOT EXISTS idx_edges_source ON edges(source)`, + `CREATE INDEX IF NOT EXISTS idx_edges_target ON edges(target)`, + `CREATE INDEX IF NOT EXISTS idx_edges_origin_status ON edges(origin, status)`, + `CREATE TABLE IF NOT EXISTS vectors ( + node_id TEXT PRIMARY KEY REFERENCES nodes(id) ON DELETE CASCADE, + dimensions INTEGER NOT NULL, + data BLOB NOT NULL, + updated_at_ns INTEGER NOT NULL +) WITHOUT ROWID`, + } + for _, statement := range statements { + if _, err := s.db.ExecContext(ctx, statement); err != nil { + return fmt.Errorf("initialize graph sqlite schema: %w", err) + } + } + var current int + err := s.db.QueryRowContext(ctx, `SELECT CAST(value AS INTEGER) FROM graph_meta WHERE key='schema_version'`).Scan(¤t) + if errors.Is(err, sql.ErrNoRows) { + _, err = s.db.ExecContext(ctx, `INSERT INTO graph_meta(key,value) VALUES('schema_version', ?)`, schemaVersion) + return err + } + if err != nil { + return err + } + if current != schemaVersion { + return fmt.Errorf("unsupported graph database schema %d (expected %d)", current, schemaVersion) + } + return nil +} + +func (s *Store) load(ctx context.Context) error { + if err := s.loadNodes(ctx); err != nil { + return err + } + if err := s.loadEdges(ctx); err != nil { + return err + } + if err := s.loadVectors(ctx); err != nil { + return err + } + var version uint64 + if err := s.db.QueryRowContext(ctx, `SELECT CAST(value AS INTEGER) FROM graph_meta WHERE key='graph_version'`).Scan(&version); err != nil && !errors.Is(err, sql.ErrNoRows) { + return err + } + var embeddingModel string + if err := s.db.QueryRowContext(ctx, `SELECT value FROM graph_meta WHERE key='embedding_model'`).Scan(&embeddingModel); err != nil && !errors.Is(err, sql.ErrNoRows) { + return err + } + var embeddingDigest string + if err := s.db.QueryRowContext(ctx, `SELECT value FROM graph_meta WHERE key='embedding_digest'`).Scan(&embeddingDigest); err != nil && !errors.Is(err, sql.ErrNoRows) { + return err + } + s.version = version + s.persistedVersion = version + s.embeddingModel = embeddingModel + s.embeddingDigest = embeddingDigest + return nil +} + +func (s *Store) loadNodes(ctx context.Context) error { + rows, err := s.db.QueryContext(ctx, `SELECT id,kind,label,summary,status,origin,external_id,uri,categories_json,keywords_json,metadata_json,weight,x,y,z,updated_at_ns FROM nodes`) + if err != nil { + return err + } + defer rows.Close() + for rows.Next() { + var n model.Node + var categories, keywords, metadata string + var updated int64 + if err := rows.Scan(&n.ID, &n.Kind, &n.Label, &n.Summary, &n.Status, &n.Origin, &n.ExternalID, &n.URI, &categories, &keywords, &metadata, &n.Weight, &n.X, &n.Y, &n.Z, &updated); err != nil { + return err + } + if err := decodeJSON(categories, &n.Categories); err != nil { + return fmt.Errorf("decode node %s categories: %w", n.ID, err) + } + if err := decodeJSON(keywords, &n.Keywords); err != nil { + return fmt.Errorf("decode node %s keywords: %w", n.ID, err) + } + if err := decodeJSON(metadata, &n.Metadata); err != nil { + return fmt.Errorf("decode node %s metadata: %w", n.ID, err) + } + if n.Metadata == nil { + n.Metadata = map[string]any{} + } + n.UpdatedAt = time.Unix(0, updated).UTC() + s.nodes[n.ID] = n + } + return rows.Err() +} + +func (s *Store) loadEdges(ctx context.Context) error { + rows, err := s.db.QueryContext(ctx, `SELECT id,source,target,type,origin,status,confidence,weight,explanation,evidence_json,metadata_json,created_at_ns,updated_at_ns FROM edges`) + if err != nil { + return err + } + defer rows.Close() + for rows.Next() { + var e model.Edge + var evidence, metadata string + var created, updated int64 + if err := rows.Scan(&e.ID, &e.Source, &e.Target, &e.Type, &e.Origin, &e.Status, &e.Confidence, &e.Weight, &e.Explanation, &evidence, &metadata, &created, &updated); err != nil { + return err + } + if err := decodeJSON(evidence, &e.Evidence); err != nil { + return fmt.Errorf("decode edge %s evidence: %w", e.ID, err) + } + if err := decodeJSON(metadata, &e.Metadata); err != nil { + return fmt.Errorf("decode edge %s metadata: %w", e.ID, err) + } + if e.Metadata == nil { + e.Metadata = map[string]any{} + } + e.CreatedAt = time.Unix(0, created).UTC() + e.UpdatedAt = time.Unix(0, updated).UTC() + s.edges[e.ID] = e + } + return rows.Err() +} + +func (s *Store) loadVectors(ctx context.Context) error { + rows, err := s.db.QueryContext(ctx, `SELECT node_id,dimensions,data FROM vectors`) + if err != nil { + return err + } + defer rows.Close() + for rows.Next() { + var id string + var dimensions int + var data []byte + if err := rows.Scan(&id, &dimensions, &data); err != nil { + return err + } + v, err := decodeVector(data, dimensions) + if err != nil { + return fmt.Errorf("decode vector %s: %w", id, err) + } + s.vectors[id] = v + } + return rows.Err() +} + +func decodeJSON(raw string, target any) error { + if strings.TrimSpace(raw) == "" { + return nil + } + dec := json.NewDecoder(strings.NewReader(raw)) + dec.UseNumber() + return dec.Decode(target) +} + +func encodeJSON(value any, empty string) (string, error) { + if value == nil { + return empty, nil + } + b, err := json.Marshal(value) + if err != nil { + return "", err + } + return string(b), nil +} + +func encodeVector(v []float32) []byte { + out := make([]byte, len(v)*4) + for i, value := range v { + binary.LittleEndian.PutUint32(out[i*4:], math.Float32bits(value)) + } + return out +} + +func decodeVector(data []byte, dimensions int) ([]float32, error) { + if dimensions < 0 || len(data) != dimensions*4 { + return nil, fmt.Errorf("invalid float32 vector blob: dimensions=%d bytes=%d", dimensions, len(data)) + } + out := make([]float32, dimensions) + for i := range out { + out[i] = math.Float32frombits(binary.LittleEndian.Uint32(data[i*4:])) + } + return out, nil +} + +// PersistVersion writes only records changed since the previous successful +// flush. Vectors are encoded as little-endian float32 BLOBs. +func (s *Store) PersistVersion() (uint64, error) { + return s.PersistVersionContext(context.Background()) +} + +func (s *Store) PersistVersionContext(ctx context.Context) (uint64, error) { + type nodeChange struct { + gen uint64 + node model.Node + } + type edgeChange struct { + gen uint64 + edge model.Edge + } + type vectorChange struct { + gen uint64 + vector []float32 + } + + s.mu.RLock() + version := s.version + persistedVersion := s.persistedVersion + embeddingModel := s.embeddingModel + embeddingDigest := s.embeddingDigest + embeddingMetaGeneration := s.embeddingMetaGeneration + nodes := make(map[string]nodeChange, len(s.dirtyNodes)) + for id, gen := range s.dirtyNodes { + if n, ok := s.nodes[id]; ok { + nodes[id] = nodeChange{gen: gen, node: n} + } + } + edges := make(map[string]edgeChange, len(s.dirtyEdges)) + for id, gen := range s.dirtyEdges { + if e, ok := s.edges[id]; ok { + edges[id] = edgeChange{gen: gen, edge: e} + } + } + vectors := make(map[string]vectorChange, len(s.dirtyVectors)) + for id, gen := range s.dirtyVectors { + if v, ok := s.vectors[id]; ok { + vectors[id] = vectorChange{gen: gen, vector: append([]float32(nil), v...)} + } + } + deletedNodes := cloneGenerations(s.deletedNodes) + deletedEdges := cloneGenerations(s.deletedEdges) + deletedVectors := cloneGenerations(s.deletedVectors) + s.mu.RUnlock() + + if len(nodes)+len(edges)+len(vectors)+len(deletedNodes)+len(deletedEdges)+len(deletedVectors) == 0 && embeddingMetaGeneration == 0 && version == persistedVersion { + return version, nil + } + + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return 0, err + } + defer tx.Rollback() + + for id := range deletedEdges { + if _, err := tx.ExecContext(ctx, `DELETE FROM edges WHERE id=?`, id); err != nil { + return 0, err + } + } + for id := range deletedNodes { + if _, err := tx.ExecContext(ctx, `DELETE FROM nodes WHERE id=?`, id); err != nil { + return 0, err + } + } + for id := range deletedVectors { + if _, err := tx.ExecContext(ctx, `DELETE FROM vectors WHERE node_id=?`, id); err != nil { + return 0, err + } + } + + if len(nodes) > 0 { + statement, err := tx.PrepareContext(ctx, nodeUpsertSQL) + if err != nil { + return 0, err + } + for _, change := range nodes { + if err := upsertNode(ctx, statement, change.node); err != nil { + _ = statement.Close() + return 0, err + } + } + if err := statement.Close(); err != nil { + return 0, err + } + } + if len(edges) > 0 { + statement, err := tx.PrepareContext(ctx, edgeUpsertSQL) + if err != nil { + return 0, err + } + for _, change := range edges { + if err := upsertEdge(ctx, statement, change.edge); err != nil { + _ = statement.Close() + return 0, err + } + } + if err := statement.Close(); err != nil { + return 0, err + } + } + if len(vectors) > 0 { + statement, err := tx.PrepareContext(ctx, vectorUpsertSQL) + if err != nil { + return 0, err + } + now := time.Now().UTC().UnixNano() + for id, change := range vectors { + if _, err := statement.ExecContext(ctx, id, len(change.vector), encodeVector(change.vector), now); err != nil { + _ = statement.Close() + return 0, err + } + } + if err := statement.Close(); err != nil { + return 0, err + } + } + if embeddingMetaGeneration > 0 { + if _, err := tx.ExecContext(ctx, `INSERT INTO graph_meta(key,value) VALUES('embedding_model',?) ON CONFLICT(key) DO UPDATE SET value=excluded.value`, embeddingModel); err != nil { + return 0, err + } + if _, err := tx.ExecContext(ctx, `INSERT INTO graph_meta(key,value) VALUES('embedding_digest',?) ON CONFLICT(key) DO UPDATE SET value=excluded.value`, embeddingDigest); err != nil { + return 0, err + } + } + if _, err := tx.ExecContext(ctx, `INSERT INTO graph_meta(key,value) VALUES('graph_version',?) ON CONFLICT(key) DO UPDATE SET value=excluded.value`, version); err != nil { + return 0, err + } + if err := tx.Commit(); err != nil { + return 0, err + } + + s.mu.Lock() + for id, change := range nodes { + if s.dirtyNodes[id] == change.gen { + delete(s.dirtyNodes, id) + } + } + for id, change := range edges { + if s.dirtyEdges[id] == change.gen { + delete(s.dirtyEdges, id) + } + } + for id, change := range vectors { + if s.dirtyVectors[id] == change.gen { + delete(s.dirtyVectors, id) + } + } + clearDeletedIfSame(s.deletedNodes, deletedNodes) + clearDeletedIfSame(s.deletedEdges, deletedEdges) + clearDeletedIfSame(s.deletedVectors, deletedVectors) + if s.embeddingMetaGeneration == embeddingMetaGeneration { + s.embeddingMetaGeneration = 0 + } + if version > s.persistedVersion { + s.persistedVersion = version + } + s.mu.Unlock() + + // Keep the WAL bounded. PASSIVE never blocks readers and leaves busy pages + // for the next scheduled flush. + _, _ = s.db.ExecContext(ctx, `PRAGMA wal_checkpoint(PASSIVE)`) + return version, nil +} + +func upsertNode(ctx context.Context, statement *sql.Stmt, n model.Node) error { + categories, err := encodeJSON(n.Categories, "[]") + if err != nil { + return err + } + keywords, err := encodeJSON(n.Keywords, "[]") + if err != nil { + return err + } + metadata, err := encodeJSON(n.Metadata, "{}") + if err != nil { + return err + } + _, err = statement.ExecContext(ctx, n.ID, n.Kind, n.Label, n.Summary, n.Status, n.Origin, n.ExternalID, n.URI, categories, keywords, metadata, n.Weight, n.X, n.Y, n.Z, n.UpdatedAt.UnixNano()) + return err +} + +func upsertEdge(ctx context.Context, statement *sql.Stmt, e model.Edge) error { + evidence, err := encodeJSON(e.Evidence, "[]") + if err != nil { + return err + } + metadata, err := encodeJSON(e.Metadata, "{}") + if err != nil { + return err + } + _, err = statement.ExecContext(ctx, e.ID, e.Source, e.Target, e.Type, e.Origin, e.Status, e.Confidence, e.Weight, e.Explanation, evidence, metadata, e.CreatedAt.UnixNano(), e.UpdatedAt.UnixNano()) + return err +} + +func cloneGenerations(in map[string]uint64) map[string]uint64 { + out := make(map[string]uint64, len(in)) + for id, gen := range in { + out[id] = gen + } + return out +} + +func clearDeletedIfSame(current, snapshot map[string]uint64) { + for id, gen := range snapshot { + if current[id] == gen { + delete(current, id) + } + } +} + +func (s *Store) Persist() error { + _, err := s.PersistVersion() + return err +} + +func (s *Store) Checkpoint(ctx context.Context, truncate bool) error { + mode := "PASSIVE" + if truncate { + mode = "TRUNCATE" + } + _, err := s.db.ExecContext(ctx, `PRAGMA wal_checkpoint(`+mode+`)`) + return err +} + +func (s *Store) Close() error { + if s.db == nil { + return nil + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + _, persistErr := s.PersistVersionContext(ctx) + checkpointErr := s.Checkpoint(ctx, true) + closeErr := s.db.Close() + return errors.Join(persistErr, checkpointErr, closeErr) +} + +func (s *Store) StorageStatus() StorageStatus { + s.mu.RLock() + var vectorBytes int64 + for _, vector := range s.vectors { + vectorBytes += int64(len(vector) * 4) + } + status := StorageStatus{ + Backend: "sqlite", Path: s.dbPath, SchemaVersion: schemaVersion, JournalMode: s.journalMode, + NodeRows: len(s.nodes), EdgeRows: len(s.edges), VectorRows: len(s.vectors), VectorBytes: vectorBytes, + PersistedVersion: s.persistedVersion, + EmbeddingModel: s.embeddingModel, + EmbeddingDigest: s.embeddingDigest, + PendingNodes: len(s.dirtyNodes), PendingEdges: len(s.dirtyEdges), PendingVectors: len(s.dirtyVectors), + PendingDeletions: len(s.deletedNodes) + len(s.deletedEdges) + len(s.deletedVectors), + PendingMetadata: s.embeddingMetaGeneration > 0, + } + s.mu.RUnlock() + if info, err := os.Stat(s.dbPath); err == nil { + status.DatabaseBytes = info.Size() + } + if info, err := os.Stat(s.dbPath + "-wal"); err == nil { + status.WALBytes = info.Size() + } + _, err := os.Stat(filepath.Join(filepath.Dir(s.dbPath), "graph-state.json")) + status.LegacyJSONPresent = err == nil + return status +} + +// Export creates a compact, transactionally consistent SQLite copy. It can be +// moved to another installation with the same knowledge-base identifiers. +func (s *Store) Export(ctx context.Context, destination string) error { + if _, err := s.PersistVersionContext(ctx); err != nil { + return err + } + if err := os.MkdirAll(filepath.Dir(destination), 0o750); err != nil { + return err + } + if err := os.Remove(destination); err != nil && !errors.Is(err, os.ErrNotExist) { + return err + } + quoted := strings.ReplaceAll(destination, "'", "''") + if _, err := s.db.ExecContext(ctx, `VACUUM INTO '`+quoted+`'`); err != nil { + return fmt.Errorf("export graph database: %w", err) + } + return os.Chmod(destination, 0o640) +} diff --git a/internal/graph/sqlite_backend_test.go b/internal/graph/sqlite_backend_test.go new file mode 100644 index 0000000..c15c128 --- /dev/null +++ b/internal/graph/sqlite_backend_test.go @@ -0,0 +1,203 @@ +package graph + +import ( + "context" + "math" + "os" + "path/filepath" + "testing" + "time" + + "github.com/local/glpi-neural-brain/internal/model" +) + +func TestSQLiteRoundTripUsesBinaryFloat32Vectors(t *testing.T) { + dir := t.TempDir() + s, err := Open(dir) + if err != nil { + t.Fatal(err) + } + s.UpsertNode(model.Node{ID: "a", Kind: "knowledge", Label: "A", Origin: "test", Categories: []string{"Netzwerk"}, Metadata: map[string]any{"rank": 2}}) + s.UpsertNode(model.Node{ID: "b", Kind: "knowledge", Label: "B", Origin: "test"}) + s.UpsertEdge(model.Edge{Source: "a", Target: "b", Type: "related_to", Origin: "test", Confidence: .8}) + s.SetVector("a", []float64{.125, -.5, .875}) + if _, err := s.PersistVersionContext(context.Background()); err != nil { + t.Fatal(err) + } + status := s.StorageStatus() + if status.Backend != "sqlite" || status.VectorBytes != 12 || status.PendingNodes != 0 || status.PendingVectors != 0 { + t.Fatalf("unexpected storage status: %+v", status) + } + if _, err := os.Stat(filepath.Join(dir, "graph-state.json")); !os.IsNotExist(err) { + t.Fatalf("legacy JSON snapshot must not be created: %v", err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + + reopened, err := Open(dir) + if err != nil { + t.Fatal(err) + } + defer reopened.Close() + if got := reopened.Snapshot(); len(got.Nodes) != 2 || len(got.Edges) != 1 { + t.Fatalf("unexpected round-trip graph: nodes=%d edges=%d", len(got.Nodes), len(got.Edges)) + } + vector, ok := reopened.Vector("a") + if !ok || len(vector) != 3 { + t.Fatalf("vector missing after round trip: %v %v", ok, vector) + } + want := []float64{.125, -.5, .875} + for i := range want { + if math.Abs(vector[i]-want[i]) > 1e-6 { + t.Fatalf("vector[%d]=%f want %f", i, vector[i], want[i]) + } + } +} + +func TestSQLitePersistsOnlyDirtyRows(t *testing.T) { + s, err := Open(t.TempDir()) + if err != nil { + t.Fatal(err) + } + defer s.Close() + s.UpsertNode(model.Node{ID: "a", Kind: "knowledge", Label: "A", Origin: "test"}) + s.UpsertNode(model.Node{ID: "b", Kind: "knowledge", Label: "B", Origin: "test"}) + s.SetVector("a", []float64{1, 0}) + if _, err := s.PersistVersion(); err != nil { + t.Fatal(err) + } + if s.Dirty() { + t.Fatal("store remained dirty after successful flush") + } + + n, _ := s.GetNode("a") + n.Label = "A2" + s.UpsertNode(n) + status := s.StorageStatus() + if status.PendingNodes != 1 || status.PendingEdges != 0 || status.PendingVectors != 0 { + t.Fatalf("expected one dirty node, got %+v", status) + } + if _, err := s.PersistVersion(); err != nil { + t.Fatal(err) + } + if status = s.StorageStatus(); status.PendingNodes != 0 || s.Dirty() { + t.Fatalf("dirty rows were not cleared: %+v", status) + } +} + +func TestSQLiteExportCreatesPortableDatabase(t *testing.T) { + s, err := Open(t.TempDir()) + if err != nil { + t.Fatal(err) + } + defer s.Close() + s.UpsertNode(model.Node{ID: "portable", Kind: "knowledge", Label: "Portable", Origin: "test"}) + s.SetVector("portable", []float64{1, 2, 3, 4}) + exportDir := t.TempDir() + destination := filepath.Join(exportDir, "graph.db") + if err := s.Export(context.Background(), destination); err != nil { + t.Fatal(err) + } + copyStore, err := Open(exportDir) + if err != nil { + t.Fatal(err) + } + defer copyStore.Close() + if _, ok := copyStore.GetNode("portable"); !ok { + t.Fatal("exported database did not contain graph node") + } + if vector, ok := copyStore.Vector("portable"); !ok || len(vector) != 4 { + t.Fatalf("exported database did not contain vector: %v %v", ok, vector) + } +} + +func TestReplaceOriginsLeavesUnchangedKnowledgeClean(t *testing.T) { + s, err := Open(t.TempDir()) + if err != nil { + t.Fatal(err) + } + defer s.Close() + updated := time.Date(2026, 8, 4, 12, 0, 0, 0, time.UTC) + nodes := []model.Node{ + {ID: "a", Kind: "knowledge", Label: "A", Summary: "gleich", Origin: "kb", UpdatedAt: updated}, + {ID: "b", Kind: "knowledge", Label: "B", Summary: "gleich", Origin: "kb", UpdatedAt: updated}, + } + edges := []model.Edge{{Source: "a", Target: "b", Type: "related_to", Origin: "kb", Weight: 1}} + s.ReplaceOrigins([]string{"kb"}, nodes, edges) + s.SetVector("a", []float64{1, 0, 0}) + if _, err := s.PersistVersion(); err != nil { + t.Fatal(err) + } + version := s.Version() + + s.ReplaceOrigins([]string{"kb"}, nodes, edges) + if s.Dirty() { + t.Fatalf("unchanged KB scan marked SQLite rows dirty: %+v", s.StorageStatus()) + } + if s.Version() != version { + t.Fatalf("unchanged KB scan changed graph version: got %d want %d", s.Version(), version) + } + if _, ok := s.Vector("a"); !ok { + t.Fatal("unchanged KB scan discarded embedding") + } +} + +func TestReplaceOriginsInvalidatesOnlyChangedEmbedding(t *testing.T) { + s, err := Open(t.TempDir()) + if err != nil { + t.Fatal(err) + } + defer s.Close() + updated := time.Date(2026, 8, 4, 12, 0, 0, 0, time.UTC) + nodes := []model.Node{{ID: "a", Kind: "knowledge", Label: "A", Summary: "alt", Origin: "kb", UpdatedAt: updated}} + s.ReplaceOrigins([]string{"kb"}, nodes, nil) + s.SetVector("a", []float64{1, 0, 0}) + if _, err := s.PersistVersion(); err != nil { + t.Fatal(err) + } + + nodes[0].Summary = "neu" + nodes[0].UpdatedAt = updated.Add(time.Minute) + s.ReplaceOrigins([]string{"kb"}, nodes, nil) + if _, ok := s.Vector("a"); ok { + t.Fatal("changed embedding input retained stale vector") + } + status := s.StorageStatus() + if status.PendingNodes != 1 || status.PendingDeletions != 1 { + t.Fatalf("expected changed node and vector deletion, got %+v", status) + } +} + +func TestEmbeddingModelIdentityInvalidatesImportedVectors(t *testing.T) { + dir := t.TempDir() + s, err := Open(dir) + if err != nil { + t.Fatal(err) + } + s.ConfigureEmbeddingIdentity("embeddinggemma", "sha256:d1") + s.UpsertNode(model.Node{ID: "a", Kind: "knowledge", Label: "A", Origin: "test"}) + s.SetVector("a", []float64{1, 2, 3}) + if _, err := s.PersistVersion(); err != nil { + t.Fatal(err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + + reopened, err := Open(dir) + if err != nil { + t.Fatal(err) + } + defer reopened.Close() + if removed := reopened.ConfigureEmbeddingIdentity("embeddinggemma", "sha256:d2"); removed != 1 { + t.Fatalf("removed vectors=%d want 1", removed) + } + if _, ok := reopened.Vector("a"); ok { + t.Fatal("vector from a different embedding model remained available") + } + status := reopened.StorageStatus() + if status.EmbeddingModel != "embeddinggemma" || status.EmbeddingDigest != "sha256:d2" || !status.PendingMetadata || status.PendingDeletions != 1 { + t.Fatalf("unexpected model-change status: %+v", status) + } +} diff --git a/internal/graph/store.go b/internal/graph/store.go index bd9de08..12e9ed7 100644 --- a/internal/graph/store.go +++ b/internal/graph/store.go @@ -2,12 +2,10 @@ package graph import ( "crypto/sha256" + "database/sql" "encoding/hex" "encoding/json" - "errors" "math" - "os" - "path/filepath" "sort" "strings" "sync" @@ -17,47 +15,29 @@ import ( ) type Store struct { - mu sync.RWMutex - nodes map[string]model.Node - edges map[string]model.Edge - vectors map[string][]float64 - version uint64 - path string - pairCursor int + mu sync.RWMutex + nodes map[string]model.Node + edges map[string]model.Edge + vectors map[string][]float32 + version uint64 + persistedVersion uint64 + pairCursor int + embeddingModel string + embeddingDigest string + embeddingMetaGeneration uint64 + + db *sql.DB + dbPath string + journalMode string + + dirtyNodes map[string]uint64 + dirtyEdges map[string]uint64 + dirtyVectors map[string]uint64 + deletedNodes map[string]uint64 + deletedEdges map[string]uint64 + deletedVectors map[string]uint64 } -type diskState struct { - Version uint64 `json:"version"` - Nodes []model.Node `json:"nodes"` - Edges []model.Edge `json:"edges"` - Vectors map[string][]float64 `json:"vectors,omitempty"` -} - -func Open(dir string) (*Store, error) { - s := &Store{nodes: map[string]model.Node{}, edges: map[string]model.Edge{}, vectors: map[string][]float64{}, path: filepath.Join(dir, "graph-state.json")} - b, err := os.ReadFile(s.path) - if errors.Is(err, os.ErrNotExist) { - return s, nil - } - if err != nil { - return nil, err - } - var d diskState - if err = json.Unmarshal(b, &d); err != nil { - return nil, err - } - s.version = d.Version - for _, n := range d.Nodes { - s.nodes[n.ID] = n - } - for _, e := range d.Edges { - s.edges[e.ID] = e - } - if d.Vectors != nil { - s.vectors = d.Vectors - } - return s, nil -} func ID(parts ...string) string { h := sha256.Sum256([]byte(strings.Join(parts, "\x00"))) return hex.EncodeToString(h[:12]) @@ -65,6 +45,84 @@ func ID(parts ...string) string { func EdgeID(source, target, typ, origin string) string { return ID("edge", source, target, typ, origin) } + +func (s *Store) markNodeDirtyLocked(id string) { + delete(s.deletedNodes, id) + s.dirtyNodes[id] = s.version +} +func (s *Store) markEdgeDirtyLocked(id string) { + delete(s.deletedEdges, id) + s.dirtyEdges[id] = s.version +} +func (s *Store) markVectorDirtyLocked(id string) { + delete(s.deletedVectors, id) + s.dirtyVectors[id] = s.version +} +func (s *Store) markNodeDeletedLocked(id string) { + delete(s.dirtyNodes, id) + s.deletedNodes[id] = s.version +} +func (s *Store) markEdgeDeletedLocked(id string) { + delete(s.dirtyEdges, id) + s.deletedEdges[id] = s.version +} +func (s *Store) markVectorDeletedLocked(id string) { + delete(s.dirtyVectors, id) + s.deletedVectors[id] = s.version +} + +func (s *Store) Dirty() bool { + s.mu.RLock() + defer s.mu.RUnlock() + return len(s.dirtyNodes)+len(s.dirtyEdges)+len(s.dirtyVectors)+len(s.deletedNodes)+len(s.deletedEdges)+len(s.deletedVectors) > 0 || s.embeddingMetaGeneration > 0 || s.version != s.persistedVersion +} + +// ConfigureEmbeddingModel records the model name before Ollama health data is +// available. A stored digest remains intact when the model name is unchanged. +func (s *Store) ConfigureEmbeddingModel(modelName string) int { + return s.ConfigureEmbeddingIdentity(modelName, "") +} + +// ConfigureEmbeddingIdentity records the model and, once known, its Ollama +// digest. Importing a database built with another model identity automatically +// invalidates all vectors while retaining nodes and edges for selective +// relearning. +func (s *Store) ConfigureEmbeddingIdentity(modelName, digest string) int { + modelName = strings.TrimSpace(modelName) + digest = strings.TrimSpace(digest) + if modelName == "" { + return 0 + } + s.mu.Lock() + defer s.mu.Unlock() + + modelChanged := s.embeddingModel != "" && s.embeddingModel != modelName + digestChanged := digest != "" && s.embeddingDigest != "" && s.embeddingDigest != digest + metadataChanged := s.embeddingModel != modelName || (digest != "" && s.embeddingDigest != digest) + if !metadataChanged { + return 0 + } + + s.version++ + removed := 0 + if modelChanged || digestChanged { + for id := range s.vectors { + delete(s.vectors, id) + s.deletedVectors[id] = s.version + delete(s.dirtyVectors, id) + removed++ + } + } + s.embeddingModel = modelName + if modelChanged { + s.embeddingDigest = digest + } else if digest != "" { + s.embeddingDigest = digest + } + s.embeddingMetaGeneration = s.version + return removed +} + func (s *Store) UpsertNode(n model.Node) { s.mu.Lock() defer s.mu.Unlock() @@ -79,6 +137,7 @@ func (s *Store) UpsertNode(n model.Node) { } s.nodes[n.ID] = n s.version++ + s.markNodeDirtyLocked(n.ID) } func (s *Store) UpsertEdge(e model.Edge) { s.mu.Lock() @@ -100,6 +159,7 @@ func (s *Store) UpsertEdge(e model.Edge) { } s.edges[e.ID] = e s.version++ + s.markEdgeDirtyLocked(e.ID) } func (s *Store) HasEdgeBetween(a, b string) bool { s.mu.RLock() @@ -130,18 +190,30 @@ func (s *Store) LookupExternal(id string) (model.Node, bool) { func (s *Store) SetVector(id string, v []float64) { s.mu.Lock() defer s.mu.Unlock() + converted := make([]float32, len(v)) + for i, value := range v { + converted[i] = float32(value) + } old, ok := s.vectors[id] - if ok && floatSlicesEqual(old, v) { + if ok && float32SlicesEqual(old, converted) { return } - s.vectors[id] = append([]float64(nil), v...) + s.vectors[id] = converted s.version++ + s.markVectorDirtyLocked(id) } func (s *Store) Vector(id string) ([]float64, bool) { s.mu.RLock() defer s.mu.RUnlock() v, ok := s.vectors[id] - return append([]float64(nil), v...), ok + if !ok { + return nil, false + } + out := make([]float64, len(v)) + for i, value := range v { + out[i] = float64(value) + } + return out, true } func (s *Store) ClearVectorsByDimension(dim int) int { @@ -152,11 +224,10 @@ func (s *Store) ClearVectorsByDimension(dim int) int { if len(v) == dim { delete(s.vectors, id) removed++ + s.version++ + s.markVectorDeletedLocked(id) } } - if removed > 0 { - s.version++ - } return removed } @@ -194,78 +265,235 @@ func (s *Store) KnowledgeNodes() []model.Node { return out } func (s *Store) ReplaceOrigins(origins []string, nodes []model.Node, edges []model.Edge) { - set := map[string]bool{} - for _, o := range origins { - set[o] = true + originSet := make(map[string]struct{}, len(origins)) + for _, origin := range origins { + originSet[origin] = struct{}{} } + + now := time.Now().UTC() + incomingNodes := make(map[string]model.Node, len(nodes)) + for _, node := range nodes { + if node.Weight == 0 { + node.Weight = 1 + } + normalizeNodeCollections(&node) + incomingNodes[node.ID] = node + } + incomingEdges := make(map[string]model.Edge, len(edges)) + for _, edge := range edges { + if edge.ID == "" { + edge.ID = EdgeID(edge.Source, edge.Target, edge.Type, edge.Origin) + } + if edge.Weight == 0 { + edge.Weight = 1 + } + normalizeEdgeCollections(&edge) + incomingEdges[edge.ID] = edge + } + s.mu.Lock() defer s.mu.Unlock() - oldVectors := make(map[string][]float64) - oldFingerprints := make(map[string]string) - for id, n := range s.nodes { - if set[n.Origin] { - if v, ok := s.vectors[id]; ok { - oldVectors[id] = append([]float64(nil), v...) - oldFingerprints[id] = n.Label + "\x00" + n.Summary + "\x00" + strings.Join(n.Categories, "\x00") + "\x00" + strings.Join(n.Keywords, "\x00") - } - delete(s.nodes, id) - delete(s.vectors, id) - } - } - for id, e := range s.edges { - if set[e.Origin] { - delete(s.edges, id) - } - } - now := time.Now().UTC() - for _, n := range nodes { - if n.UpdatedAt.IsZero() { - n.UpdatedAt = now - } - if n.Weight == 0 { - n.Weight = 1 - } - if n.X == 0 && n.Y == 0 && n.Z == 0 { - n.X, n.Y, n.Z = position(n.ID, n.Categories) - } - s.nodes[n.ID] = n - fingerprint := n.Label + "\x00" + n.Summary + "\x00" + strings.Join(n.Categories, "\x00") + "\x00" + strings.Join(n.Keywords, "\x00") - if fingerprint == oldFingerprints[n.ID] { - if v, ok := oldVectors[n.ID]; ok { - s.vectors[n.ID] = append([]float64(nil), v...) - } - } - } - for _, e := range edges { - if e.ID == "" { - e.ID = EdgeID(e.Source, e.Target, e.Type, e.Origin) - } - if e.CreatedAt.IsZero() { - e.CreatedAt = now - } - e.UpdatedAt = now - if e.Weight == 0 { - e.Weight = 1 - } - s.edges[e.ID] = e - } - for id, e := range s.edges { - if _, ok := s.nodes[e.Source]; !ok { - delete(s.edges, id) + + // Remove records from managed origins that disappeared from the source. + for id, old := range s.nodes { + if _, managed := originSet[old.Origin]; !managed { continue } - if _, ok := s.nodes[e.Target]; !ok { - delete(s.edges, id) + if _, present := incomingNodes[id]; present { + continue + } + delete(s.nodes, id) + if _, hadVector := s.vectors[id]; hadVector { + delete(s.vectors, id) + s.version++ + s.markVectorDeletedLocked(id) + } + s.version++ + s.markNodeDeletedLocked(id) + } + for id, old := range s.edges { + if _, managed := originSet[old.Origin]; !managed { + continue + } + if _, present := incomingEdges[id]; present { + continue + } + delete(s.edges, id) + s.version++ + s.markEdgeDeletedLocked(id) + } + + // Reconcile nodes instead of deleting and recreating every row on each scan. + // This is what makes the scheduled SQLite flush truly incremental for an + // unchanged knowledge base. + for id, incoming := range incomingNodes { + old, existed := s.nodes[id] + if incoming.UpdatedAt.IsZero() { + if existed && !old.UpdatedAt.IsZero() { + incoming.UpdatedAt = old.UpdatedAt + } else { + incoming.UpdatedAt = now + } + } + if incoming.X == 0 && incoming.Y == 0 && incoming.Z == 0 { + if existed && (old.X != 0 || old.Y != 0 || old.Z != 0) { + incoming.X, incoming.Y, incoming.Z = old.X, old.Y, old.Z + } else { + incoming.X, incoming.Y, incoming.Z = position(incoming.ID, incoming.Categories) + } + } + + oldEmbeddingFingerprint := "" + if existed { + oldEmbeddingFingerprint = embeddingFingerprint(old) + } + if existed && nodesEquivalent(old, incoming) { + continue + } + + s.nodes[id] = incoming + s.version++ + s.markNodeDirtyLocked(id) + + // Embeddings depend on text/categories/keywords, not on display + // coordinates or unrelated metadata. Only invalidate a vector when its + // actual embedding input changed. + if existed && oldEmbeddingFingerprint != embeddingFingerprint(incoming) { + if _, hadVector := s.vectors[id]; hadVector { + delete(s.vectors, id) + s.version++ + s.markVectorDeletedLocked(id) + } + } + } + + // Reconcile deterministic source edges. Timestamps are retained for an + // unchanged edge so periodic scans do not produce needless writes. + for id, incoming := range incomingEdges { + old, existed := s.edges[id] + if existed && edgesEquivalentIgnoringTimestamps(old, incoming) { + continue + } + if incoming.CreatedAt.IsZero() { + if existed && !old.CreatedAt.IsZero() { + incoming.CreatedAt = old.CreatedAt + } else { + incoming.CreatedAt = now + } + } + incoming.UpdatedAt = now + s.edges[id] = incoming + s.version++ + s.markEdgeDirtyLocked(id) + } + + // Remove any remaining edge whose endpoint no longer exists. This includes + // AI-derived edges that referred to a source note removed from the KB. + for id, edge := range s.edges { + if _, ok := s.nodes[edge.Source]; !ok { + delete(s.edges, id) + s.version++ + s.markEdgeDeletedLocked(id) + continue + } + if _, ok := s.nodes[edge.Target]; !ok { + delete(s.edges, id) + s.version++ + s.markEdgeDeletedLocked(id) } } - s.version++ } + +func embeddingFingerprint(node model.Node) string { + return node.Label + "\x00" + node.Summary + "\x00" + strings.Join(node.Categories, "\x00") + "\x00" + strings.Join(node.Keywords, "\x00") +} + +func normalizeNodeCollections(node *model.Node) { + if node.Categories == nil { + node.Categories = []string{} + } + if node.Keywords == nil { + node.Keywords = []string{} + } + if node.Metadata == nil { + node.Metadata = map[string]any{} + } +} + +func normalizeEdgeCollections(edge *model.Edge) { + if edge.Evidence == nil { + edge.Evidence = []model.Evidence{} + } + if edge.Metadata == nil { + edge.Metadata = map[string]any{} + } +} + +func nodesEquivalent(a, b model.Node) bool { + if a.ID != b.ID || a.Kind != b.Kind || a.Label != b.Label || + a.Summary != b.Summary || a.Status != b.Status || a.Origin != b.Origin || + a.ExternalID != b.ExternalID || a.URI != b.URI || a.Weight != b.Weight || + a.X != b.X || a.Y != b.Y || a.Z != b.Z { + return false + } + return JSONEquivalent(a.Categories, b.Categories) && + JSONEquivalent(a.Keywords, b.Keywords) && + JSONEquivalent(a.Metadata, b.Metadata) +} + +func edgesEquivalentIgnoringTimestamps(a, b model.Edge) bool { + if a.ID != b.ID || a.Source != b.Source || a.Target != b.Target || + a.Type != b.Type || a.Origin != b.Origin || a.Status != b.Status || + a.Confidence != b.Confidence || a.Weight != b.Weight || + a.Explanation != b.Explanation { + return false + } + return JSONEquivalent(a.Evidence, b.Evidence) && JSONEquivalent(a.Metadata, b.Metadata) +} + +func JSONEquivalent(a, b any) bool { + left, leftErr := json.Marshal(a) + right, rightErr := json.Marshal(b) + return leftErr == nil && rightErr == nil && string(left) == string(right) +} + func (s *Store) Version() uint64 { s.mu.RLock() defer s.mu.RUnlock() return s.version } +func (s *Store) Counts() (nodes, edges int, version uint64) { + s.mu.RLock() + defer s.mu.RUnlock() + nodes = len(s.nodes) + for _, edge := range s.edges { + if edge.Status != "rejected" { + edges++ + } + } + return nodes, edges, s.version +} + +func (s *Store) IdleNode(seed int64) (model.Node, bool) { + s.mu.RLock() + defer s.mu.RUnlock() + if len(s.nodes) == 0 { + return model.Node{}, false + } + index := int(seed % int64(len(s.nodes))) + if index < 0 { + index = -index + } + for _, node := range s.nodes { + if index == 0 { + return node, true + } + index-- + } + return model.Node{}, false +} + func (s *Store) Snapshot() model.Snapshot { s.mu.RLock() defer s.mu.RUnlock() @@ -284,37 +512,6 @@ func (s *Store) Snapshot() model.Snapshot { sort.Slice(e, func(i, j int) bool { return e[i].ID < e[j].ID }) return model.Snapshot{Version: s.version, Nodes: n, Edges: e, UpdatedAt: time.Now().UTC()} } -func (s *Store) Persist() error { - _, err := s.PersistVersion() - return err -} - -func (s *Store) PersistVersion() (uint64, error) { - s.mu.RLock() - d := diskState{Version: s.version, Vectors: map[string][]float64{}} - for _, n := range s.nodes { - d.Nodes = append(d.Nodes, n) - } - for _, e := range s.edges { - d.Edges = append(d.Edges, e) - } - for k, v := range s.vectors { - d.Vectors[k] = append([]float64(nil), v...) - } - s.mu.RUnlock() - b, err := json.Marshal(d) - if err != nil { - return 0, err - } - tmp := s.path + ".tmp" - if err = os.WriteFile(tmp, b, 0o640); err != nil { - return 0, err - } - if err := os.Rename(tmp, s.path); err != nil { - return 0, err - } - return d.Version, nil -} func (s *Store) Similar(query []float64, limit int) []model.Hit { s.mu.RLock() defer s.mu.RUnlock() @@ -324,7 +521,7 @@ func (s *Store) Similar(query []float64, limit int) []model.Hit { if !ok || (n.Kind != "knowledge" && n.Kind != "ai-think" && n.Kind != "external") { continue } - score := cosine(query, v) + score := cosineMixed(query, v) hits = append(hits, model.Hit{NodeID: id, Label: n.Label, Score: score, Kind: n.Kind, Status: n.Status}) } sort.Slice(hits, func(i, j int) bool { return hits[i].Score > hits[j].Score }) @@ -402,7 +599,7 @@ func (s *Store) NextPairFilteredDepth(min float64, anchorLimit int, categories [ continue } comparisons++ - score := cosine(lv, rv) + score := cosine32(lv, rv) if score >= min && score > best { best = score a, b = left, right @@ -438,7 +635,7 @@ func (s *Store) BestPair(min float64) (model.Node, model.Node, float64, bool) { if edgeBetweenLocked(s.edges, nodes[i].ID, nodes[j].ID) { continue } - score := cosine(s.vectors[nodes[i].ID], s.vectors[nodes[j].ID]) + score := cosine32(s.vectors[nodes[i].ID], s.vectors[nodes[j].ID]) if score >= min && score > best { best = score a = nodes[i] @@ -523,7 +720,7 @@ func edgeBetweenLocked(edges map[string]model.Edge, a, b string) bool { } return false } -func floatSlicesEqual(a, b []float64) bool { +func float32SlicesEqual(a, b []float32) bool { if len(a) != len(b) { return false } @@ -535,15 +732,33 @@ func floatSlicesEqual(a, b []float64) bool { return true } -func cosine(a, b []float64) float64 { +func cosine32(a, b []float32) float64 { if len(a) == 0 || len(a) != len(b) { return 0 } var dot, aa, bb float64 for i := range a { - dot += a[i] * b[i] + av, bv := float64(a[i]), float64(b[i]) + dot += av * bv + aa += av * av + bb += bv * bv + } + if aa == 0 || bb == 0 { + return 0 + } + return dot / (math.Sqrt(aa) * math.Sqrt(bb)) +} + +func cosineMixed(a []float64, b []float32) float64 { + if len(a) == 0 || len(a) != len(b) { + return 0 + } + var dot, aa, bb float64 + for i := range a { + bv := float64(b[i]) + dot += a[i] * bv aa += a[i] * a[i] - bb += b[i] * b[i] + bb += bv * bv } if aa == 0 || bb == 0 { return 0 diff --git a/internal/graph/store_test.go b/internal/graph/store_test.go index 082b175..b5044a6 100644 --- a/internal/graph/store_test.go +++ b/internal/graph/store_test.go @@ -11,6 +11,7 @@ func TestReplaceOriginsPreservesUnchangedVector(t *testing.T) { if err != nil { t.Fatal(err) } + t.Cleanup(func() { _ = s.Close() }) n := model.Node{ID: "n1", Kind: "knowledge", Label: "VPN", Summary: "Gateway", Origin: "knowledge-production"} s.ReplaceOrigins([]string{"knowledge-production"}, []model.Node{n}, nil) s.SetVector("n1", []float64{1, 2, 3}) @@ -28,6 +29,7 @@ func TestReplaceOriginsPreservesUnchangedVector(t *testing.T) { func TestRejectedEdgePreventsPairReprocessingButIsHidden(t *testing.T) { s, _ := Open(t.TempDir()) + t.Cleanup(func() { _ = s.Close() }) a := model.Node{ID: "a", Kind: "knowledge", Label: "A", Origin: "knowledge-production"} b := model.Node{ID: "b", Kind: "knowledge", Label: "B", Origin: "knowledge-production"} s.UpsertNode(a) @@ -48,6 +50,7 @@ func TestRejectedEdgePreventsPairReprocessingButIsHidden(t *testing.T) { func TestNextPairUsesBoundedRotatingAnchors(t *testing.T) { s, _ := Open(t.TempDir()) + t.Cleanup(func() { _ = s.Close() }) for _, id := range []string{"a", "b", "c", "d"} { s.UpsertNode(model.Node{ID: id, Kind: "knowledge", Label: id, Origin: "knowledge-production"}) } @@ -75,6 +78,7 @@ func TestCategoryFiltersLimitEmbeddingAndThinkingCandidates(t *testing.T) { if err != nil { t.Fatal(err) } + t.Cleanup(func() { _ = s.Close() }) nodes := []model.Node{ {ID: "net-a", Kind: "knowledge", Label: "Net A", Categories: []string{"Netzwerk"}, Origin: "test"}, {ID: "net-b", Kind: "knowledge", Label: "Net B", Categories: []string{"Netzwerk"}, Origin: "test"}, diff --git a/internal/persist/coordinator.go b/internal/persist/coordinator.go index 7f46bdc..0ba4237 100644 --- a/internal/persist/coordinator.go +++ b/internal/persist/coordinator.go @@ -127,12 +127,15 @@ func (c *Coordinator) Flush(ctx context.Context, trigger string) error { currentVersion = c.graph.Version() } graphDirty := currentVersion != lastVersion + if c.graph != nil { + graphDirty = c.graph.Dirty() + } if len(files) == 0 && !graphDirty { return nil } - // Knowledge/staging/cache files are committed first. The graph snapshot is - // written afterwards so a persisted graph never points at a draft that has + // Knowledge/staging/cache files are committed first. The incremental SQLite + // transaction follows so a persisted graph never points at a draft that has // not yet reached disk. for _, f := range files { if err := writeAtomic(f.path, f.data, f.perm); err != nil { @@ -147,7 +150,7 @@ func (c *Coordinator) Flush(ctx context.Context, trigger string) error { } if graphDirty && c.graph != nil { - persistedVersion, err := c.graph.PersistVersion() + persistedVersion, err := c.graph.PersistVersionContext(ctx) if err != nil { c.recordFailure(err) return err @@ -164,7 +167,7 @@ func (c *Coordinator) Flush(ctx context.Context, trigger string) error { remaining := len(c.files) c.mu.Unlock() if c.broker != nil { - c.broker.Publish(model.Activity{Type: "persistence.flushed", Source: "brain", Phase: "storage", Message: "Graph und Wissensdateien wurden gebündelt geschrieben", Strength: .32, Metadata: map[string]any{"trigger": trigger, "files": len(files), "graph_version": currentVersion, "remaining_files": remaining}}) + c.broker.Publish(model.Activity{Type: "persistence.flushed", Source: "brain", Phase: "storage", Message: "SQLite-Graph und Wissensdateien wurden gebündelt geschrieben", Strength: .32, Metadata: map[string]any{"trigger": trigger, "files": len(files), "graph_version": currentVersion, "remaining_files": remaining}}) } return nil } @@ -184,9 +187,14 @@ func (c *Coordinator) Status() Status { c.mu.RLock() defer c.mu.RUnlock() return Status{ - Interval: c.interval.String(), - PendingFiles: len(c.files), - GraphDirty: current != c.lastPersistedVersion, + Interval: c.interval.String(), + PendingFiles: len(c.files), + GraphDirty: func() bool { + if c.graph != nil { + return c.graph.Dirty() + } + return current != c.lastPersistedVersion + }(), LastFlush: c.lastFlush, LastError: c.lastError, LastPersistedVersion: c.lastPersistedVersion, diff --git a/internal/persist/coordinator_test.go b/internal/persist/coordinator_test.go index e1a3bdf..e38a200 100644 --- a/internal/persist/coordinator_test.go +++ b/internal/persist/coordinator_test.go @@ -11,7 +11,7 @@ import ( "github.com/local/glpi-neural-brain/internal/model" ) -func TestCoordinatorBatchesFilesAndGraph(t *testing.T) { +func TestCoordinatorBatchesFilesAndSQLiteGraph(t *testing.T) { dir := t.TempDir() g, err := graph.Open(dir) if err != nil { @@ -29,8 +29,11 @@ func TestCoordinatorBatchesFilesAndGraph(t *testing.T) { if _, err := os.Stat(path); !os.IsNotExist(err) { t.Fatalf("file should remain buffered before flush: %v", err) } + if _, err := os.Stat(filepath.Join(dir, "graph.db")); err != nil { + t.Fatalf("SQLite database should be initialized at open: %v", err) + } if _, err := os.Stat(filepath.Join(dir, "graph-state.json")); !os.IsNotExist(err) { - t.Fatalf("graph should remain buffered before flush: %v", err) + t.Fatalf("legacy JSON graph must not be written: %v", err) } if err := c.Flush(context.Background(), "test"); err != nil { t.Fatal(err) @@ -42,11 +45,20 @@ func TestCoordinatorBatchesFilesAndGraph(t *testing.T) { if string(data) != "second" { t.Fatalf("unexpected file data %q", data) } - if _, err := os.Stat(filepath.Join(dir, "graph-state.json")); err != nil { - t.Fatal(err) - } status := c.Status() if status.PendingFiles != 0 || status.GraphDirty || status.SuccessfulFlushes != 1 { t.Fatalf("unexpected status: %+v", status) } + + if err := g.Close(); err != nil { + t.Fatal(err) + } + reopened, err := graph.Open(dir) + if err != nil { + t.Fatal(err) + } + defer reopened.Close() + if _, ok := reopened.GetNode("n1"); !ok { + t.Fatal("graph row was not persisted by coordinator") + } } diff --git a/internal/research/searxng.go b/internal/research/searxng.go index b0ba613..a08b61a 100644 --- a/internal/research/searxng.go +++ b/internal/research/searxng.go @@ -3,40 +3,151 @@ package research import ( "context" "encoding/json" + "errors" "fmt" + "io" + "net" "net/http" "net/url" "strings" + "sync" "time" "github.com/local/glpi-neural-brain/internal/model" ) +const maxResponseBytes = 8 << 20 + +type Diagnostic struct { + Configured bool `json:"configured"` + BaseURL string `json:"base_url,omitempty"` + OK bool `json:"ok"` + CheckedAt time.Time `json:"checked_at,omitempty"` + DurationMS int64 `json:"duration_ms,omitempty"` + HTTPStatus int `json:"http_status,omitempty"` + ContentType string `json:"content_type,omitempty"` + ResultCount int `json:"result_count,omitempty"` + ErrorKind string `json:"error_kind,omitempty"` + Error string `json:"error,omitempty"` +} + type Client struct { BaseURL string HTTP *http.Client + + mu sync.RWMutex + last Diagnostic } func New(base string) *Client { - return &Client{BaseURL: strings.TrimRight(base, "/"), HTTP: &http.Client{Timeout: 45 * time.Second}} + base = strings.TrimSpace(base) + return &Client{ + BaseURL: strings.TrimRight(base, "/"), + HTTP: &http.Client{Timeout: 45 * time.Second}, + last: Diagnostic{ + Configured: base != "", + BaseURL: displayBaseURL(base), + }, + } } + +func (c *Client) Status() Diagnostic { + if c == nil { + return Diagnostic{} + } + c.mu.RLock() + out := c.last + c.mu.RUnlock() + if out.BaseURL == "" { + out.BaseURL = displayBaseURL(c.BaseURL) + } + out.Configured = strings.TrimSpace(c.BaseURL) != "" + return out +} + func (c *Client) Search(ctx context.Context, q string, limit int) ([]model.ResearchResult, error) { - if c == nil || c.BaseURL == "" { - return nil, nil + results, _, err := c.SearchDetailed(ctx, q, limit) + return results, err +} + +func (c *Client) SearchDetailed(ctx context.Context, q string, limit int) ([]model.ResearchResult, Diagnostic, error) { + started := time.Now() + diagnostic := Diagnostic{ + Configured: c != nil && strings.TrimSpace(c.BaseURL) != "", + CheckedAt: started.UTC(), } - u := c.BaseURL + "/search?format=json&language=de-DE&q=" + url.QueryEscape(q) - req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) - if err != nil { - return nil, err + if c != nil { + diagnostic.BaseURL = displayBaseURL(c.BaseURL) } - resp, err := c.HTTP.Do(req) + finish := func(results []model.ResearchResult, err error) ([]model.ResearchResult, Diagnostic, error) { + diagnostic.DurationMS = time.Since(started).Milliseconds() + diagnostic.ResultCount = len(results) + diagnostic.OK = err == nil + if err != nil { + diagnostic.Error = err.Error() + if diagnostic.ErrorKind == "" { + diagnostic.ErrorKind = classifyError(err) + } + } + if c != nil { + c.mu.Lock() + c.last = diagnostic + c.mu.Unlock() + } + return results, diagnostic, err + } + + if c == nil || strings.TrimSpace(c.BaseURL) == "" { + diagnostic.ErrorKind = "not_configured" + return finish(nil, errors.New("SearXNG ist nicht konfiguriert: SEARXNG_URL ist leer")) + } + q = strings.TrimSpace(q) + if q == "" { + diagnostic.ErrorKind = "invalid_query" + return finish(nil, errors.New("SearXNG-Suchanfrage ist leer")) + } + + u, err := searchURL(c.BaseURL, q) if err != nil { - return nil, err + diagnostic.ErrorKind = "invalid_url" + return finish(nil, fmt.Errorf("ungültige SEARXNG_URL %q: %w", diagnostic.BaseURL, err)) + } + req, err := http.NewRequestWithContext(ctx, http.MethodGet, u.String(), nil) + if err != nil { + diagnostic.ErrorKind = "request_build" + return finish(nil, fmt.Errorf("SearXNG-Anfrage konnte nicht erstellt werden: %w", err)) + } + req.Header.Set("Accept", "application/json") + req.Header.Set("Accept-Language", "de-DE,de;q=0.9,en;q=0.7") + req.Header.Set("User-Agent", "glpi-neural-brain/1.0") + + client := c.HTTP + if client == nil { + client = &http.Client{Timeout: 45 * time.Second} + } + resp, err := client.Do(req) + if err != nil { + diagnostic.ErrorKind = classifyError(err) + return finish(nil, fmt.Errorf("SearXNG unter %s ist nicht erreichbar: %w%s", endpointForError(u), err, networkHint(u, diagnostic.ErrorKind))) } defer resp.Body.Close() - if resp.StatusCode/100 != 2 { - return nil, fmt.Errorf("searxng HTTP %d", resp.StatusCode) + diagnostic.HTTPStatus = resp.StatusCode + diagnostic.ContentType = strings.TrimSpace(resp.Header.Get("Content-Type")) + + body, readErr := io.ReadAll(io.LimitReader(resp.Body, maxResponseBytes+1)) + if readErr != nil { + diagnostic.ErrorKind = "response_read" + return finish(nil, fmt.Errorf("SearXNG-Antwort konnte nicht gelesen werden: %w", readErr)) } + if len(body) > maxResponseBytes { + diagnostic.ErrorKind = "response_too_large" + return finish(nil, fmt.Errorf("SearXNG-Antwort überschreitet %d MiB", maxResponseBytes>>20)) + } + if resp.StatusCode/100 != 2 { + diagnostic.ErrorKind = "http_status" + return finish(nil, fmt.Errorf("SearXNG lieferte HTTP %d (%s)%s%s", resp.StatusCode, http.StatusText(resp.StatusCode), bodySuffix(body), httpStatusHint(resp.StatusCode))) + } + var env struct { Results []struct { Title string `json:"title"` @@ -44,18 +155,147 @@ func (c *Client) Search(ctx context.Context, q string, limit int) ([]model.Resea Content string `json:"content"` } `json:"results"` } - if err := json.NewDecoder(resp.Body).Decode(&env); err != nil { - return nil, err + if err := json.Unmarshal(body, &env); err != nil { + diagnostic.ErrorKind = "invalid_json" + return finish(nil, fmt.Errorf("SearXNG lieferte kein gültiges JSON (Content-Type %q): %w%s; in settings.yml muss search.formats den Wert json enthalten", diagnostic.ContentType, err, bodySuffix(body))) } - out := []model.ResearchResult{} + + out := make([]model.ResearchResult, 0, len(env.Results)) for _, r := range env.Results { if strings.TrimSpace(r.URL) == "" { continue } - out = append(out, model.ResearchResult{Title: r.Title, URL: r.URL, Content: r.Content}) + out = append(out, model.ResearchResult{Title: strings.TrimSpace(r.Title), URL: strings.TrimSpace(r.URL), Content: strings.TrimSpace(r.Content)}) if limit > 0 && len(out) >= limit { break } } - return out, nil + return finish(out, nil) +} + +func searchURL(base, query string) (*url.URL, error) { + u, err := url.Parse(strings.TrimSpace(base)) + if err != nil { + return nil, err + } + if u.Scheme != "http" && u.Scheme != "https" { + return nil, fmt.Errorf("Schema muss http oder https sein") + } + if u.Host == "" { + return nil, fmt.Errorf("Host fehlt") + } + path := strings.TrimRight(u.Path, "/") + if !strings.HasSuffix(strings.ToLower(path), "/search") { + path += "/search" + } + if path == "" { + path = "/search" + } + u.Path = path + values := u.Query() + values.Set("format", "json") + values.Set("language", "de-DE") + values.Set("q", query) + u.RawQuery = values.Encode() + return u, nil +} + +func displayBaseURL(raw string) string { + u, err := url.Parse(strings.TrimSpace(raw)) + if err != nil { + return strings.TrimSpace(raw) + } + u.User = nil + u.RawQuery = "" + u.Fragment = "" + return strings.TrimRight(u.String(), "/") +} + +func endpointForError(u *url.URL) string { + if u == nil { + return "unbekannter Adresse" + } + clean := *u + clean.User = nil + clean.RawQuery = "" + clean.Fragment = "" + return clean.String() +} + +func classifyError(err error) string { + if err == nil { + return "" + } + if errors.Is(err, context.DeadlineExceeded) { + return "timeout" + } + if errors.Is(err, context.Canceled) { + return "canceled" + } + text := strings.ToLower(err.Error()) + switch { + case strings.Contains(text, "connection refused"): + return "connection_refused" + case strings.Contains(text, "no such host"): + return "dns" + case strings.Contains(text, "certificate") || strings.Contains(text, "tls"): + return "tls" + } + var dnsErr *net.DNSError + if errors.As(err, &dnsErr) { + return "dns" + } + var netErr net.Error + if errors.As(err, &netErr) { + if netErr.Timeout() { + return "timeout" + } + return "network" + } + return "request" +} + +func networkHint(u *url.URL, kind string) string { + if u == nil { + return "" + } + host := strings.ToLower(u.Hostname()) + if host == "localhost" || host == "127.0.0.1" || host == "::1" { + return "; Hinweis: Innerhalb des Brain-Containers verweist localhost auf den Brain-Container. Verwende den Docker-Service-Namen, zum Beispiel http://searxng:8080, oder host.docker.internal" + } + switch kind { + case "dns": + return "; prüfe Docker-Netzwerk und Service-Namen" + case "connection_refused", "network": + return "; am Zielhost ist auf diesem Port möglicherweise kein SearXNG-Dienst erreichbar" + case "tls": + return "; prüfe Zertifikat und HTTPS-Konfiguration" + default: + return "" + } +} + +func httpStatusHint(status int) string { + switch status { + case http.StatusForbidden: + return "; prüfe SearXNG-Limiter/Zugriffsregeln und ob das JSON-Format erlaubt ist" + case http.StatusNotFound: + return "; prüfe die SearXNG-Basis-URL und einen möglichen Reverse-Proxy-Pfad" + case http.StatusTooManyRequests: + return "; SearXNG oder ein vorgeschalteter Proxy begrenzt die Anfrage" + default: + return "" + } +} + +func bodySuffix(body []byte) string { + text := strings.Join(strings.Fields(string(body)), " ") + if text == "" { + return "" + } + const max = 320 + if len(text) > max { + text = text[:max] + "…" + } + return ": " + text } diff --git a/internal/research/searxng_test.go b/internal/research/searxng_test.go index c9840f2..54fdb76 100644 --- a/internal/research/searxng_test.go +++ b/internal/research/searxng_test.go @@ -4,20 +4,75 @@ import ( "context" "net/http" "net/http/httptest" + "strings" "testing" ) func TestSearchParsesSearXNG(t *testing.T) { srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/search" { + t.Fatalf("unexpected path %q", r.URL.Path) + } + if got := r.URL.Query().Get("format"); got != "json" { + t.Fatalf("unexpected format %q", got) + } w.Header().Set("Content-Type", "application/json") _, _ = w.Write([]byte(`{"results":[{"title":"Vendor note","url":"https://example.test/a","content":"Relevant evidence"}]}`)) })) defer srv.Close() - got, err := New(srv.URL).Search(context.Background(), "vpn", 3) + got, diagnostic, err := New(srv.URL).SearchDetailed(context.Background(), "vpn", 3) if err != nil { t.Fatal(err) } if len(got) != 1 || got[0].URL != "https://example.test/a" || got[0].Content != "Relevant evidence" { t.Fatalf("unexpected result: %#v", got) } + if !diagnostic.OK || diagnostic.HTTPStatus != http.StatusOK || diagnostic.ResultCount != 1 { + t.Fatalf("unexpected diagnostic: %+v", diagnostic) + } +} + +func TestSearchAcceptsBaseURLWithSearchPath(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/search" { + t.Fatalf("search path duplicated: %q", r.URL.Path) + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"results":[]}`)) + })) + defer srv.Close() + if _, err := New(srv.URL+"/search").Search(context.Background(), "test", 1); err != nil { + t.Fatal(err) + } +} + +func TestSearchReportsHTTPBody(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/plain") + w.WriteHeader(http.StatusForbidden) + _, _ = w.Write([]byte("JSON format is disabled")) + })) + defer srv.Close() + _, diagnostic, err := New(srv.URL).SearchDetailed(context.Background(), "test", 1) + if err == nil || !strings.Contains(err.Error(), "HTTP 403") || !strings.Contains(err.Error(), "JSON format is disabled") { + t.Fatalf("unexpected error: %v", err) + } + if diagnostic.ErrorKind != "http_status" || diagnostic.HTTPStatus != http.StatusForbidden { + t.Fatalf("unexpected diagnostic: %+v", diagnostic) + } +} + +func TestSearchExplainsInvalidJSON(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/html") + _, _ = w.Write([]byte("search page")) + })) + defer srv.Close() + _, diagnostic, err := New(srv.URL).SearchDetailed(context.Background(), "test", 1) + if err == nil || !strings.Contains(err.Error(), "settings.yml") { + t.Fatalf("unexpected error: %v", err) + } + if diagnostic.ErrorKind != "invalid_json" { + t.Fatalf("unexpected diagnostic: %+v", diagnostic) + } } diff --git a/internal/web/server.go b/internal/web/server.go index c8f9baf..c133a95 100644 --- a/internal/web/server.go +++ b/internal/web/server.go @@ -4,9 +4,12 @@ import ( "context" "embed" "encoding/json" + "fmt" "io" "io/fs" "net/http" + "os" + "path/filepath" "strings" "time" @@ -34,6 +37,8 @@ func (s *Server) Handler() http.Handler { mux.HandleFunc("GET /api/runtime-settings", s.handleGetRuntimeSettings) mux.HandleFunc("PUT /api/runtime-settings", s.handleSetRuntimeSettings) mux.HandleFunc("GET /api/categories", s.handleCategories) + mux.HandleFunc("GET /api/research/status", s.handleResearchStatus) + mux.HandleFunc("POST /api/research/test", s.handleResearchTest) mux.HandleFunc("GET /api/stream", s.Broker.ServeSSE) mux.HandleFunc("POST /api/query", s.handleQuery) mux.HandleFunc("POST /api/events", s.handleEvent) @@ -41,6 +46,7 @@ func (s *Server) Handler() http.Handler { mux.HandleFunc("POST /api/enrich", s.handleEnrich) mux.HandleFunc("POST /api/glpi-kb/sync", s.handleGLPIKBSync) mux.HandleFunc("POST /api/flush", s.handleFlush) + mux.HandleFunc("GET /api/state/export", s.handleStateExport) sub, _ := fs.Sub(assets, "static") mux.Handle("GET /", http.FileServer(http.FS(sub))) return s.headers(mux) @@ -90,6 +96,40 @@ func (s *Server) handleCategories(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, map[string]any{"categories": s.Engine.Categories()}) } +func (s *Server) handleResearchStatus(w http.ResponseWriter, r *http.Request) { + writeJSON(w, http.StatusOK, s.Engine.ResearchStatus()) +} + +func (s *Server) handleResearchTest(w http.ResponseWriter, r *http.Request) { + if !s.authorized(r) { + writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "unauthorized"}) + return + } + var request struct { + Query string `json:"query"` + Limit int `json:"limit"` + } + if r.Body != nil { + err := json.NewDecoder(io.LimitReader(r.Body, 1<<20)).Decode(&request) + if err != nil && err != io.EOF { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()}) + return + } + } + ctx, cancel := contextTimeout(r, 60*time.Second) + defer cancel() + result, err := s.Engine.TestResearch(ctx, request.Query, request.Limit) + if err != nil { + writeJSON(w, http.StatusBadGateway, map[string]any{ + "error": err.Error(), + "query": result.Query, + "diagnostic": result.Diagnostic, + }) + return + } + writeJSON(w, http.StatusOK, result) +} + func (s *Server) handleQuery(w http.ResponseWriter, r *http.Request) { if !s.authorized(r) { writeJSON(w, 401, map[string]string{"error": "unauthorized"}) @@ -155,6 +195,41 @@ func (s *Server) handleGLPIKBSync(w http.ResponseWriter, r *http.Request) { writeJSON(w, 200, s.Engine.Status()) } +func (s *Server) handleStateExport(w http.ResponseWriter, r *http.Request) { + if !s.authorized(r) { + writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "unauthorized"}) + return + } + dir, err := os.MkdirTemp("", "brain-export-*") + if err != nil { + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()}) + return + } + defer os.RemoveAll(dir) + path := filepath.Join(dir, "graph.db") + ctx, cancel := contextTimeout(r, 10*time.Minute) + defer cancel() + if err := s.Engine.ExportGraph(ctx, path); err != nil { + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()}) + return + } + file, err := os.Open(path) + if err != nil { + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()}) + return + } + defer file.Close() + info, err := file.Stat() + if err != nil { + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()}) + return + } + w.Header().Set("Content-Type", "application/vnd.sqlite3") + w.Header().Set("Content-Disposition", `attachment; filename="graph.db"`) + w.Header().Set("Content-Length", fmt.Sprintf("%d", info.Size())) + http.ServeContent(w, r, "graph.db", info.ModTime(), file) +} + func (s *Server) handleFlush(w http.ResponseWriter, r *http.Request) { if !s.authorized(r) { writeJSON(w, 401, map[string]string{"error": "unauthorized"}) diff --git a/internal/web/server_test.go b/internal/web/server_test.go index 8ff5ddc..2da46a6 100644 --- a/internal/web/server_test.go +++ b/internal/web/server_test.go @@ -13,6 +13,7 @@ import ( "github.com/local/glpi-neural-brain/internal/engine" "github.com/local/glpi-neural-brain/internal/graph" "github.com/local/glpi-neural-brain/internal/model" + "github.com/local/glpi-neural-brain/internal/research" ) func TestRuntimeSettingsAndCategoriesAPI(t *testing.T) { @@ -79,3 +80,47 @@ func TestRuntimeSettingsAndCategoriesAPI(t *testing.T) { t.Fatalf("disabled thinking should return 409, got %d: %s", res.Code, res.Body.String()) } } + +func TestResearchDiagnosticAPI(t *testing.T) { + searx := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"results":[{"title":"Btrfs documentation","url":"https://example.test/btrfs","content":"Snapshot evidence"}]}`)) + })) + defer searx.Close() + + broker := activity.New(20) + eng := &engine.Engine{ + Cfg: config.Config{ResearchEnabled: true}, + Research: research.New(searx.URL), + Broker: broker, + } + h := (&Server{Engine: eng, Broker: broker}).Handler() + + res := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodPost, "/api/research/test", strings.NewReader(`{"query":"btrfs snapshots","limit":4}`)) + req.Header.Set("Content-Type", "application/json") + h.ServeHTTP(res, req) + if res.Code != http.StatusOK { + t.Fatalf("research test returned %d: %s", res.Code, res.Body.String()) + } + var result engine.ResearchTestResult + if err := json.NewDecoder(res.Body).Decode(&result); err != nil { + t.Fatal(err) + } + if !result.OK || !result.Diagnostic.OK || len(result.Results) != 1 { + t.Fatalf("unexpected result: %+v", result) + } + + res = httptest.NewRecorder() + h.ServeHTTP(res, httptest.NewRequest(http.MethodGet, "/api/research/status", nil)) + if res.Code != http.StatusOK { + t.Fatalf("research status returned %d: %s", res.Code, res.Body.String()) + } + var status map[string]any + if err := json.NewDecoder(res.Body).Decode(&status); err != nil { + t.Fatal(err) + } + if status["ok"] != true { + t.Fatalf("unexpected status: %#v", status) + } +} diff --git a/internal/web/static/app.css b/internal/web/static/app.css index b782b27..0f45840 100644 --- a/internal/web/static/app.css +++ b/internal/web/static/app.css @@ -50,3 +50,4 @@ body.honeycomb-view .legend{opacity:.58}body.honeycomb-view .mode-status small:a .activity-item .sources{display:grid;gap:4px;margin-top:7px;padding-top:6px;border-top:1px solid rgba(93,255,189,.1)} .activity-item .sources span{display:flex;align-items:flex-start;gap:6px;color:#b7d9ca;font-size:9px;line-height:1.35} .activity-item .sources i{display:inline-flex;align-items:center;justify-content:center;flex:0 0 15px;height:15px;border-radius:50%;border:1px solid rgba(93,255,189,.24);background:rgba(93,255,189,.07);color:#86ffd0;font-style:normal;font-size:8px} +.research-settings{flex:0 0 auto}.research-status{margin:8px 0 10px;padding:9px 10px;border:1px solid rgba(133,200,255,.11);border-radius:11px;background:rgba(255,255,255,.022);color:#87a2b5;font-size:9px;line-height:1.45;word-break:break-word}.research-status.ok{border-color:rgba(93,255,189,.22);color:#91d9bc;background:rgba(93,255,189,.045)}.research-status.failed{border-color:rgba(255,105,125,.24);color:#ff9aab;background:rgba(255,105,125,.045)}.research-status strong{display:block;color:inherit;font-size:10px;margin-bottom:2px}.research-test-query{display:block}.research-test-query span{display:block;font-size:9px;color:#6d879a;margin-bottom:5px}.research-test-query input{width:100%;border:1px solid rgba(133,200,255,.14);background:rgba(2,8,17,.6);color:var(--text);border-radius:10px;padding:9px 10px;outline:0}.research-test-query input:focus{border-color:rgba(93,255,189,.4);box-shadow:0 0 0 3px rgba(93,255,189,.06)}.research-test-actions{display:flex;align-items:center;gap:9px;margin-top:8px}.research-test-actions button{border:1px solid rgba(93,255,189,.3);background:rgba(93,255,189,.08);color:var(--green);border-radius:10px;padding:8px 10px;font-size:8px;font-weight:800;letter-spacing:.12em;cursor:pointer}.research-test-actions button:disabled{opacity:.5;cursor:wait}.research-test-actions span{font-size:9px;color:#7893a6;line-height:1.35}.activity-item .error-detail{margin-top:7px;padding:7px 8px;border-left:2px solid rgba(255,105,125,.6);background:rgba(255,105,125,.055);color:#ff9aab;font-size:9px;line-height:1.45;white-space:pre-wrap;word-break:break-word} diff --git a/internal/web/static/app.js b/internal/web/static/app.js index 679998b..a73ec07 100644 --- a/internal/web/static/app.js +++ b/internal/web/static/app.js @@ -63,7 +63,12 @@ async function api(path, options = {}) { const res = await fetch(path, {headers: {'Content-Type': 'application/json', ...(options.headers || {})}, ...options}); const data = await res.json().catch(() => ({})); - if (!res.ok) throw new Error(data.error || `HTTP ${res.status}`); + if (!res.ok) { + const error = new Error(data.error || `HTTP ${res.status}`); + error.status = res.status; + error.data = data; + throw error; + } return data; } @@ -389,11 +394,47 @@ syncRuntimeControls(); } renderAutonomyStatus(status); + renderResearchStatus(status.searxng || {configured: Boolean(status.research_enabled)}); } catch { renderAutonomyStatus({ollama_ok: false, auto_enrich: false, enrich_error: 'Status nicht erreichbar'}); } } + function renderResearchStatus(status = {}) { + const box = $('researchStatus'); + const badge = $('researchStatusBadge'); + if (!box || !badge) return; + const configured = Boolean(status.configured); + const checked = Boolean(status.checked_at); + const ok = Boolean(status.ok); + const endpoint = status.base_url || 'SEARXNG_URL nicht gesetzt'; + box.className = `research-status ${checked ? (ok ? 'ok' : 'failed') : ''}`; + if (!configured) { + badge.textContent = 'deaktiviert'; + box.innerHTML = `Nicht konfiguriert${escapeHTML(status.error || 'BRAIN_RESEARCH_ENABLED und SEARXNG_URL prüfen.')}`; + return; + } + if (status.error_kind === 'disabled') { + badge.textContent = 'deaktiviert'; + box.innerHTML = `${escapeHTML(endpoint)}${escapeHTML(status.error || 'BRAIN_RESEARCH_ENABLED ist deaktiviert.')}`; + return; + } + if (!checked) { + badge.textContent = 'bereit'; + box.innerHTML = `${escapeHTML(endpoint)}Noch keine echte JSON-Suche ausgeführt.`; + return; + } + if (ok) { + badge.textContent = 'erreichbar'; + const count = Number(status.result_count || 0); + box.innerHTML = `${escapeHTML(endpoint)}HTTP ${Number(status.http_status || 200)} · ${count} Treffer · ${Number(status.duration_ms || 0).toLocaleString('de-DE')} ms`; + return; + } + badge.textContent = 'Fehler'; + const kind = status.error_kind ? `${escapeHTML(String(status.error_kind))} · ` : ''; + box.innerHTML = `${escapeHTML(endpoint)}${kind}${escapeHTML(status.error || 'SearXNG-Test fehlgeschlagen')}`; + } + function renderAutonomyStatus(status) { const panel = $('autonomyStatus'); const button = $('enrichNow'); @@ -1958,12 +1999,12 @@ } addLog(evt); if (evt.type === 'graph.updated') loadGraph(); - if (evt.type?.startsWith('think.') || evt.type?.startsWith('article.')) loadStatus(); + if (evt.type?.startsWith('think.') || evt.type?.startsWith('article.') || evt.type?.startsWith('research.')) loadStatus(); } function shouldLog(evt) { if (!evt || evt.type === 'brain.idle' || evt.type === 'node.activated' || evt.type === 'edges.traversed') return false; - const important = new Set(['scan.started', 'graph.updated', 'embedding.batch', 'query.started', 'query.completed', 'think.queued', 'think.cycle.started', 'think.cycle.completed', 'think.cycle.failed', 'think.no_candidate', 'think.started', 'think.relation.created', 'think.rejected', 'think.failed', 'think.paused', 'research.started', 'research.results', 'research.ingested', 'research.failed', 'article.plan.started', 'article.plan.skipped', 'article.research.started', 'article.research.results', 'article.research.ingested', 'article.research.failed', 'article.draft.started', 'article.draft.rejected', 'article.created', 'article.duplicate', 'article.skipped', 'article.failed', 'agent.run', 'glpi.kb.synced', 'glpi.kb.failed', 'persistence.flushed', 'persistence.failed']); + const important = new Set(['scan.started', 'graph.updated', 'embedding.batch', 'query.started', 'query.completed', 'think.queued', 'think.cycle.started', 'think.cycle.completed', 'think.cycle.failed', 'think.no_candidate', 'think.started', 'think.relation.created', 'think.rejected', 'think.failed', 'think.paused', 'research.started', 'research.results', 'research.ingested', 'research.failed', 'research.test.started', 'research.test.results', 'research.test.failed', 'article.plan.started', 'article.plan.skipped', 'article.research.started', 'article.research.results', 'article.research.ingested', 'article.research.failed', 'article.draft.started', 'article.draft.rejected', 'article.created', 'article.duplicate', 'article.skipped', 'article.failed', 'agent.run', 'glpi.kb.synced', 'glpi.kb.failed', 'persistence.flushed', 'persistence.failed']); if (!important.has(evt.type) && !(evt.source === 'agent' || evt.source === 'knowledgebase' || evt.source === 'external' || evt.query)) return false; const fingerprint = `${evt.type}|${evt.message || ''}|${evt.query || ''}|${evt.source || ''}`; const last = state.lastLogFingerprint.get(fingerprint) || 0; @@ -1999,6 +2040,9 @@ 'research.results': 'SearXNG-Ergebnisse empfangen', 'research.ingested': 'Webquellen im Graph verknüpft', 'research.failed': 'Relationsrecherche fehlgeschlagen', + 'research.test.started': 'SearXNG-Test gestartet', + 'research.test.results': 'SearXNG-Test erfolgreich', + 'research.test.failed': 'SearXNG-Test fehlgeschlagen', 'article.plan.started': 'Artikelmehrwert wird geprüft', 'article.plan.skipped': 'Artikelsynthese übersprungen', 'article.research.started': 'Artikelrecherche gestartet', @@ -2034,6 +2078,9 @@ if (evt.metadata?.research_result_count) meta.push(`${Number(evt.metadata.research_result_count)} Webquellen`); if (evt.metadata?.result_count !== undefined) meta.push(`${Number(evt.metadata.result_count)} SearXNG-Treffer`); if (Array.isArray(evt.metadata?.result_domains) && evt.metadata.result_domains.length) meta.push(evt.metadata.result_domains.slice(0, 3).join(' · ')); + if (evt.metadata?.searxng_http_status) meta.push(`HTTP ${Number(evt.metadata.searxng_http_status)}`); + if (evt.metadata?.error_kind) meta.push(String(evt.metadata.error_kind)); + if (evt.metadata?.searxng_content_type) meta.push(String(evt.metadata.searxng_content_type).split(';')[0]); if (evt.metadata?.trigger) meta.push(evt.metadata.trigger === 'manual' ? 'manuell' : 'automatisch'); if (evt.metadata?.batch_size) meta.push(`${Number(evt.metadata.batch_size)} Schritte`); if (evt.metadata?.checked !== undefined) meta.push(`${Number(evt.metadata.checked)} geprüft`); @@ -2064,7 +2111,7 @@ message = `${evt.source === 'agent' ? 'Agent' : evt.source === 'knowledgebase' ? 'Knowledgebase' : 'Brain'} verarbeitet eine Anfrage.`; } const sources = Array.isArray(evt.metadata?.result_titles) ? evt.metadata.result_titles.map(String).slice(0, 4) : []; - return {time, title, message, meta, query: eventQuery, regions, sources, cls: evt.type?.includes('think') || evt.type?.startsWith('article.') ? (evt.type?.includes('research') ? 'research' : 'think') : evt.type?.includes('research') ? 'research' : evt.type === 'graph.updated' || evt.type === 'scan.started' || evt.type?.startsWith('glpi.kb') || evt.type?.startsWith('persistence.') ? 'graph' : evt.source === 'agent' ? 'agent' : ''}; + return {time, title, message, meta, query: eventQuery, regions, sources, error: evt.metadata?.error ? String(evt.metadata.error) : '', endpoint: evt.metadata?.searxng_base_url ? String(evt.metadata.searxng_base_url) : '', cls: evt.type?.includes('think') || evt.type?.startsWith('article.') ? (evt.type?.includes('research') ? 'research' : 'think') : evt.type?.includes('research') ? 'research' : evt.type === 'graph.updated' || evt.type === 'scan.started' || evt.type?.startsWith('glpi.kb') || evt.type?.startsWith('persistence.') ? 'graph' : evt.source === 'agent' ? 'agent' : ''}; } function addLog(evt) { @@ -2072,7 +2119,7 @@ const out = formatEvent(evt); const item = document.createElement('div'); item.className = 'activity-item ' + out.cls; - item.innerHTML = `${escapeHTML(out.title)}
${escapeHTML(out.message)}
${out.regions.length ? `${escapeHTML(out.message)}
${out.regions.length ? `Unbegrenzt · Kategorie-Filter und LOD bestimmen die Renderlast.
+