Update 6 SQLite
All checks were successful
release-tag / release-image (push) Successful in 2m30s

This commit is contained in:
2026-08-05 05:23:43 +02:00
parent 7e36c2f3f0
commit b4388d1ebc
40 changed files with 2515 additions and 230 deletions

View File

@@ -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

View File

@@ -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

View File

@@ -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`.

View File

@@ -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.

View File

@@ -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.

View File

@@ -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.

View File

@@ -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

View File

@@ -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

View File

@@ -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.

View File

@@ -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`.

View File

@@ -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.

42
SQLITE-STORAGE.md Normal file
View File

@@ -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.

View File

@@ -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`.

16
VALIDATION-SQLITE.md Normal file
View File

@@ -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.

View File

@@ -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)
}
}

0
data/.gitkeep Normal file
View File

View File

@@ -1 +0,0 @@
{"version":4,"nodes":null,"edges":null}

BIN
data/graph.db Normal file

Binary file not shown.

View File

@@ -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
}

View File

@@ -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

View File

@@ -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}

15
go.mod
View File

@@ -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
)

47
go.sum Normal file
View File

@@ -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=

View File

@@ -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"

View File

@@ -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

View File

@@ -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
}

View File

@@ -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)
}
}

View File

@@ -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(&current)
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)
}

View File

@@ -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)
}
}

View File

@@ -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

View File

@@ -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"},

View File

@@ -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,

View File

@@ -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")
}
}

View File

@@ -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
}

View File

@@ -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("<html>search page</html>"))
}))
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)
}
}

View File

@@ -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"})

View File

@@ -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)
}
}

View File

@@ -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}

View File

@@ -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 = `<strong>Nicht konfiguriert</strong>${escapeHTML(status.error || 'BRAIN_RESEARCH_ENABLED und SEARXNG_URL prüfen.')}`;
return;
}
if (status.error_kind === 'disabled') {
badge.textContent = 'deaktiviert';
box.innerHTML = `<strong>${escapeHTML(endpoint)}</strong>${escapeHTML(status.error || 'BRAIN_RESEARCH_ENABLED ist deaktiviert.')}`;
return;
}
if (!checked) {
badge.textContent = 'bereit';
box.innerHTML = `<strong>${escapeHTML(endpoint)}</strong>Noch keine echte JSON-Suche ausgeführt.`;
return;
}
if (ok) {
badge.textContent = 'erreichbar';
const count = Number(status.result_count || 0);
box.innerHTML = `<strong>${escapeHTML(endpoint)}</strong>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 = `<strong>${escapeHTML(endpoint)}</strong>${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 = `<b>${escapeHTML(out.title)}</b><time>${out.time}</time><p>${escapeHTML(out.message)}</p>${out.regions.length ? `<div class="region">Cortex: ${out.regions.map(escapeHTML).join(' · ')}</div>` : ''}${out.meta.length ? `<div class="meta">${out.meta.map(v => `<span>${escapeHTML(v)}</span>`).join('')}</div>` : ''}${out.sources?.length ? `<div class="sources">${out.sources.map((source, index) => `<span><i>${index + 1}</i>${escapeHTML(source)}</span>`).join('')}</div>` : ''}${out.query ? `<div class="query">${escapeHTML(out.query)}</div>` : ''}`;
item.innerHTML = `<b>${escapeHTML(out.title)}</b><time>${out.time}</time><p>${escapeHTML(out.message)}</p>${out.regions.length ? `<div class="region">Cortex: ${out.regions.map(escapeHTML).join(' · ')}</div>` : ''}${out.meta.length ? `<div class="meta">${out.meta.map(v => `<span>${escapeHTML(v)}</span>`).join('')}</div>` : ''}${out.sources?.length ? `<div class="sources">${out.sources.map((source, index) => `<span><i>${index + 1}</i>${escapeHTML(source)}</span>`).join('')}</div>` : ''}${out.query ? `<div class="query">${escapeHTML(out.query)}</div>` : ''}${out.error ? `<div class="error-detail">${out.endpoint ? `${escapeHTML(out.endpoint)}\n` : ''}${escapeHTML(out.error)}</div>` : ''}`;
const log = $('activityLog');
log.prepend(item);
while (log.children.length > 18) log.removeChild(log.lastChild);
@@ -2123,6 +2170,25 @@
}
});
$('testResearch')?.addEventListener('click', async e => {
const button = e.currentTarget;
const feedback = $('researchTestFeedback');
const query = String($('researchTestQuery')?.value || '').trim();
button.disabled = true;
if (feedback) feedback.textContent = 'Test läuft …';
try {
const result = await api('/api/research/test', {method: 'POST', body: JSON.stringify({query, limit: 4})});
renderResearchStatus(result.diagnostic || {});
if (feedback) feedback.textContent = `${Number(result.diagnostic?.result_count || result.results?.length || 0)} Treffer empfangen`;
} catch (err) {
renderResearchStatus(err.data?.diagnostic || {configured: true, ok: false, checked_at: new Date().toISOString(), error: err.message});
if (feedback) feedback.textContent = err.message;
} finally {
button.disabled = false;
await loadStatus();
}
});
$('openSettings').addEventListener('click', openSettingsPanel);
$('closeSettings').addEventListener('click', closeSettingsPanel);
$('settingsBackdrop').addEventListener('click', closeSettingsPanel);

View File

@@ -108,6 +108,19 @@
<p id="displayLimitHint" class="setting-hint">Unbegrenzt · Kategorie-Filter und LOD bestimmen die Renderlast.</p>
</section>
<section class="settings-section research-settings">
<div class="settings-section-title"><h2>SearXNG-Diagnose</h2><span id="researchStatusBadge">nicht geprüft</span></div>
<div id="researchStatus" class="research-status">Konfiguration wird geladen …</div>
<label class="research-test-query" for="researchTestQuery">
<span>Direkte Testanfrage</span>
<input id="researchTestQuery" type="text" value="Btrfs Snapshots und ZFS History Unterschiede Timeline" autocomplete="off">
</label>
<div class="research-test-actions">
<button id="testResearch" type="button">SEARXNG TESTEN</button>
<span id="researchTestFeedback" aria-live="polite"></span>
</div>
</section>
<section class="settings-section category-settings">
<div class="settings-section-title"><h2>Kategorie-Filter</h2><span>Leer = alle Kategorien</span></div>
<label class="filter-search"><span>Kategorien durchsuchen</span><input id="categorySearch" type="search" placeholder="z. B. GLPI, Netzwerk, Ollama"></label>